- 数据库
- 后端
- 流处理
【免费下载链接】EventStore
KurrentDB is a database that's engineered for modern software applications and event-driven architectures. Its event-native design simplifies data modeling and preserves data integrity while the integrated streaming engine solves distributed messaging challenges and ensures data consistency.
KurrentDB Connectors 是运行在服务器端、把 KurrentDB 事件流接入外部系统(消息队列、数据库、HTTP 端点等)的组件,每个连接器内部由订阅、过滤/转换、Sink 或 Source 等环节组成。本篇指南围绕连接器公开的指标体系展开:如何从/metrics端点获取连接器运行指标、每一类指标(连接器、Sink、Consumer、Producer、Processor)对应的时间序列与类型含义,以及如何结合源码理解这些指标的产生位置,为数据管道建立可观测的监控与告警体系。
指标从何而来:连接器的可观测性设计
连接器的数据处理链路在 intro.md 中有清晰描述:连接器使用 catch-up 订阅接收事件,经过过滤(Filter)与转换(Transform)后,通过 Sink 推送到外部系统。这条链路中的每一个关键环节——订阅消费、消息生产、Sink 写入、数据转换——都对应一组可观测指标,这正是 metrics.md 将指标按Usage / Sink / Consumer / Producer / Processor五类组织的原因。
从源码结构看,指标的产生点与连接器组件一一对应:
- Sink 相关指标由各 Sink 实现负责记录,例如 SerilogSink.cs 中每次写入调用
SinkMetrics.TrackWrite(context.ConnectorId, MetricsLabel),其中MetricsLabel返回"serilog"等 Sink 类型标识; - 连接器生命周期指标由工厂记录,SystemConnectorsFactory.cs 在创建/关闭 Sink 或 Source 连接器时分别调用
ConnectorMetrics.TrackSinkConnectorCreated/Closed与TrackSourceConnectorCreated/Closed,并在转换失败、规约失败时调用SinkMetrics.TrackTransformError、TrackReduceError; - Consumer 与 Producer 指标由底层消息系统(Surge 框架)的拦截器产生,SystemConsumer.cs 与 SystemProducer.cs 分别注册
ConsumerMetrics与ProducerMetrics拦截器。
连接器控制面的指标则由 ControlPlaneMetrics.cs 统一创建,Meter 名为EventStore.Connectors.ControlPlane(见 DiagnosticsName.cs)。
指标如何暴露:/metrics 端点与 metricsconfig.json
KurrentDB 以 Prometheus 文本格式 中对指标体系与metricsconfig.json的说明。
连接器指标是否进入/metrics输出,取决于metricsconfig.json中的Meters配置。在仓库的 metricsconfig.json 中,与连接器相关的 Meter 已默认启用:
"Meters": [ "KurrentDB.Core", "KurrentDB.Projections.Core", "Kurrent", "Kurrent.Connectors", "Kurrent.Connectors.Sinks", "KurrentDB.SecondaryIndexes" ]其中Kurrent.Connectors对应连接器生命周期指标(如kurrent_connector_active_total),Kurrent.Connectors.Sinks对应 Sink 写入类指标。同时,该文件顶部的ExpectedScrapeIntervalSeconds(必须为0、1、5、10或15的倍数)决定了RecentMax类型指标统计窗口的大小,间接影响连接器指标中延迟/耗时类指标的尖峰捕获能力,相关原理见 metrics.md。
此外,KurrentDB 还支持通过 OpenTelemetry Protocol (OTLP) 主动向指定端点导出指标,连接器指标同样可以随之上报。
连接器级指标(Usage):活跃连接器计数
Usage 类别关注连接器整体的"存活性",即当前正在运行的连接器数量:
| 时间序列 | 类型 | 描述 |
|---|---|---|
kurrent_connector_active_total | Gauge | 当前活跃的数据连接器数量 |
该指标为 Gauge,含义是"此刻有多少个连接器处于运行状态"。结合源码可以更精确地理解其语义:在 SystemConnectorsFactoryTests.cs 中有一条名为counts_only_disposed_connector_as_closed的测试,它通过MeterListener监听Kurrent.ConnectorsMeter 下的kurrent_connector_active_total,断言只有连接器被DisposeAsync之后才会计为"关闭"。这条测试印证了指标的生命周期语义:
- 连接器创建(
CreateConnector)→ 指标 +1; - 连接器正常关闭(
DisposeAsync,即使 Sink 关闭失败,见 SystemConnectorsFactory.cs 中的注释)→ 指标 -1; - 指标带有
connector_id标签,可用于区分具体是哪个连接器。
因此,当kurrent_connector_active_total与预期不符(如连接器启动失败却未计数、或关闭后仍计数)时,应优先检查连接器的生命周期管理。
Sink 指标:写入链路的核心观测点
Sink 指标回答"数据是否成功写入了目标系统、写得多快、出错了没有"这三个问题:
| 时间序列 | 类型 | 描述 |
|---|---|---|
kurrent_sink_written_total_records | Histogram | 成功写入 Sink 的记录总数 |
kurrent_sink_errors_total | Counter | Sink 操作期间遇到的总错误数 |
kurrent_sink_transform_duration_s | Histogram | 写入 Sink 之前的数据转换耗时(秒) |
kurrent_sink_write_latency_s | Histogram | 事件创建到 Sink 写入确认之间的时间(秒) |
需要特别说明两点:
kurrent_sink_written_total_records虽名为 Histogram,但语义是"累计总数",其_sum/_count可用于计算写入速率(records/s),结合kurrent_sink_write_latency_s的分布可以判断写入瓶颈是发生在网络传输还是目标系统本身。kurrent_sink_transform_duration_s衡量的是数据转换环节的耗时。连接器的转换功能是把 JavaScript 编写的转换函数以 base64 编码配置在连接器上(参考 features.md),转换会直接影响 Sink 写入前的处理耗时。从源码看,转换器通过JintRecordTransformer执行,并在出错时回调SinkMetrics.TrackTransformError记录错误(见 SystemConnectorsFactory.cs);SQL 类 Sink 还额外通过JintSqlReducer做字段规约,其错误回调为SinkMetrics.TrackReduceError。
kurrent_sink_errors_total是 Counter(只增不减),若要判断"当前是否在持续出错",应观察其在一段时间内的增量(rate),而不是瞬时值。
Consumer 指标:订阅消费侧
Consumer 指标描述连接器从 KurrentDB 订阅并消费消息的情况,由 Consumer 拦截器产生:
| 时间序列 | 类型 | 描述 |
|---|---|---|
messaging_kurrent_consumer_message_count | Counter | 从消息系统消费的消息总数 |
messaging_kurrent_consumer_commit_latency_s | Histogram | 收到记录与其位置被提交之间的时间(秒) |
messaging_kurrent_consumer_lag | Gauge | 最新消息与最后一条已消费消息之间的差值 |
其中最有监控价值的是messaging_kurrent_consumer_lag:它直接反映"订阅是否跟上了写入速度"。若 lag 持续增长,说明消费处理(下游 Sink 写入)速度跟不上事件产生速度,是整个数据管道的核心瓶颈信号。
messaging_kurrent_consumer_commit_latency_s则与检查点(checkpoint)机制相关。连接器会周期性地把已成功处理的最后事件位置写入$connectors/{connector-id}/checkpoints系统流(见 features.md),提交延迟升高意味着检查点提交缓慢,可能拖累故障恢复时的"续传"精度。从源码看,消费路径在 SystemConsumer.cs 中通过CheckpointController提交位置,并注册了ConsumerMetrics拦截器来记录消费侧指标。
Producer 指标:Source 连接器的消息生产侧
Source 连接器(如 kafka.md 中描述的外部消息源)从外部系统拉取消息并生产到 KurrentDB,Producer 指标描述这一侧的运行情况:
| 时间序列 | 类型 | 描述 |
|---|---|---|
messaging_kurrent_producer_queue_length | Gauge | 生产者队列中等待发送的消息数量 |
messaging_kurrent_producer_message_count | Counter | 成功生产到消息系统的消息总数 |
messaging_kurrent_producer_produce_duration_s | Histogram | 向消息系统生产消息所花费的时间(秒) |
messaging_kurrent_producer_queue_length是观察生产背压(backpressure)的直接指标:队列长度持续走高说明下游(KurrentDB 写入)处理不过来,此时事件在连接器内部排队。生产路径由 SystemProducer.cs 实现,它同样注册了ProducerMetrics拦截器,并在 SystemConnectorsFactory.cs 中通过SystemProducer.Builder为每个 Source 连接器创建 Producer 实例。
Processor 指标:消息处理环节
Processor 指标记录连接器内部消息处理环节的错误情况:
| 时间序列 | 类型 | 描述 |
|---|---|---|
messaging_kurrent_processor_error_count | Counter | 消息处理期间遇到的总错误数 |
该指标与kurrent_sink_errors_total的关注点不同:后者是"写入外部系统时出错",前者是"连接器内部处理消息时出错"(例如反序列化失败、过滤/转换异常等)。两者配合使用,可以区分故障是出在下游目标系统还是连接器自身处理逻辑。从源码结构看,连接器采用"Processor(处理器)+ Interceptor(拦截器)"架构,SystemProcessor.cs 与 SystemConnectorsFactory.cs 展示了处理器如何串联客户端、状态存储、Schema 注册表、过滤器和 Sink 代理,处理环节的异常会反映到 Processor 指标上。
指标类型速查:如何读懂 Gauge、Counter 与 Histogram
连接器指标使用了三种常见类型,其语义定义可参考 KurrentDB 官方指标文档 metrics.md 中的"Common types"一节(对应 Prometheus 指标类型说明):
- Gauge:当前值,可升可降。用于描述"此刻的状态",如
kurrent_connector_active_total、messaging_kurrent_consumer_lag、messaging_kurrent_producer_queue_length。监控这类指标应关注其绝对值与变化趋势。 - Counter:累计计数,只增不减。用于描述"到目前为止总共发生了多少次",如
kurrent_sink_errors_total、messaging_kurrent_consumer_message_count、messaging_kurrent_processor_error_count。监控这类指标应使用rate()/increase()计算增量。 - Histogram:观测值分布。用于描述延迟/耗时的分位数与分布,如
kurrent_sink_transform_duration_s、kurrent_sink_write_latency_s、messaging_kurrent_consumer_commit_latency_s、messaging_kurrent_producer_produce_duration_s。查询时通常结合histogram_quantile计算 p50/p95/p99。
此外,KurrentDB 还存在一种特殊的RecentMax类型(连接器指标未直接使用,但ExpectedScrapeIntervalSeconds配置影响所有基于 RecentMax 的指标窗口):它记录"一组最近测量值中的最大值",用于捕获两次抓取之间可能漏掉的尖峰,详见 metrics.md。
监控与告警实践建议
基于以上指标体系,可以搭建一套覆盖"连接器存活 → 消费进度 → 写入健康"三层的数据管道监控:
- 存活与数量:对
kurrent_connector_active_total设置告警,当其低于预期连接器数量时,说明有连接器异常退出; - 消费进度:对
messaging_kurrent_consumer_lag设置阈值告警(如持续 5 分钟大于 N),lag 持续增长通常意味着下游处理能力不足; - 写入健康:对
rate(kurrent_sink_errors_total[5m])与rate(messaging_kurrent_processor_error_count[5m])设置非零告警,并结合kurrent_sink_write_latency_s的 p99 判断是否需要对 Sink 目标系统扩容; - 转换性能:若配置了转换/规约函数(参考 features.md),关注
kurrent_sink_transform_duration_s的分布,异常偏高说明转换函数(base64 编码的 JavaScript)存在性能问题; - 背压:对 Source 类连接器关注
messaging_kurrent_producer_queue_length,持续增长说明消息生产快于 KurrentDB 写入。
所有指标均可在 KurrentDB 的/metrics端点直接抓取,无需额外部署 exporter。官方还提供了 Cluster Summary 与 miscellaneous panels 两套 Grafana 面板,可作为连接器指标看板设计的起点(见 metrics.md 开头部分)。
- 数据库
- 后端
- 流处理
【免费下载链接】EventStore
KurrentDB is a database that's engineered for modern software applications and event-driven architectures. Its event-native design simplifies data modeling and preserves data integrity while the integrated streaming engine solves distributed messaging challenges and ensures data consistency.
相关推荐
Bindu Agent 健康检查与监控指标指南:/health 与 /metrics 端点深度解析
Bindu Agent 健康检查与监控指标指南:/health 与 /metrics 端点深度解析 健康检查与监控是任何 AI Agent 上生产环境的第一步。
人工智能AI Agent认证鉴权后端RPC框架Frigate 监控指标(Metrics)完全指南:用 Prometheus 与 Grafana 观测 NVR 健康与性能
Frigate 监控指标(Metrics)完全指南:用 Prometheus 与 Grafana 观测 NVR 健康与性能 Frigate 是一套面向 IP 摄
人工智能计算机视觉音视频Apache Doris数据监控:性能指标与健康检查
Apache Doris数据监控:性能指标与健康检查 你是否还在为分布式数据库的性能问题头疼?当数据量达到TB级别,查询延迟突然飙升,却找不到问题根源?Apac
数据库OLAP大数据数据仓库分布式数据库实时分析列式数据库
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考