资讯动态

SeaTunnel Prometheus Sink 数据接收器实战:Remote Write 协议、参数配置与检查点刷新机制

发布时间:2026/9/18 21:25:26 来源:尧图企业网站定制
SeaTunnel Prometheus Sink 数据接收器实战Remote Write 协议、参数配置与检查点刷新机制【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel本文基于 Apache SeaTunnel 仓库中connector-prometheus连接器模块及 Prometheus Sink 官方文档系统讲解 SeaTunnel 如何将上游数据通过 Prometheus Remote Write 协议Snappy 压缩 HTTP POST写入 Prometheus 及其兼容后端如 VictoriaMetrics覆盖全部接收器参数、字段类型契约、三种刷新时机batch 触发 / 引擎定时 / 检查点以及失败重放语义并深入源码与单元测试揭示底层实现。概述把 SeaTunnel 数据流接到 PrometheusPrometheus 数据接收器Sink是 SeaTunnel 连接器家族中负责出站可观测性数据管道的一环它把上游任意来源Kafka、JDBC、FakeSource 等的 SeaTunnel 行数据转换为 Prometheus Remote Write 协议要求的采样点用 Snappy 压缩后通过 HTTPPOST请求写入 Prometheus 兼容的 remote write 地址。典型的应用场景包括把业务事件流、日志指标同步进 Prometheus 时序库做告警与监控将 SeaTunnel 处理后的指标写入 VictoriaMetrics、Mimir、Thanos 等兼容 Remote Write 的后端作为跨系统指标投递管道替代手写 exporter 或 agent。读完本文你将掌握三个核心字段标签、数值、时间戳如何映射到采样点、全部接收器选项的含义与默认值、sink.flush.interval定时刷新的引擎级实现原理以及检查点刷新失败时的数据重放安全边界。引擎支持与主要特性根据 Prometheus.md 与连接器实现该 Sink 支持以下运行引擎SparkFlinkSeaTunnel Zeta特性清单如下对照 连接器 V2 特性说明特性是否支持精准一次Exactly Once否[ ]变更数据捕获CDC否[ ]多表写入Multi Table Sink是[x]定时刷新Timer Flush是[x]仅 SeaTunnel Zeta 引擎其中多表写入由SupportMultiTableSink接口支撑见 PrometheusSink.java定时刷新是引擎级能力详见下文定时刷新小节并非连接器自身维护线程。工作原理从 SeaTunnel 行到 Prometheus 采样点Prometheus Sink 从上游数据中取出3 个字段来组成一个 Prometheus 采样点sample字段含义建议类型key_labelPrometheus 标签字段mapstring, stringkey_value指标数值字段doublekey_timestamp可选的时间戳字段见下文类型矩阵这三个字段名通过配置指定key_label/key_value/key_timestamp对应源码中 PrometheusSerializer.java 构造的三个提取器label/value/timestamp extractor将每行数据序列化为 Point 对象包含metric: MapString,String、value: Double、timestamp: Long。随后 PrometheusWriter.createRemoteWriteRequest() 会把 Point 列表构建成 protobuf 格式的WriteRequestRemote.WriteRequest时间序列TimeSeries 标签Label 采样点Sample再经 Snappy 压缩后 POST 出去。请求地址形如http://prometheus:9090/api/v1/writehttp://victoria-metrics:8428/api/v1/write注意如果样本时间戳太早目标服务可能会按自身的保留策略retention或 remote write 规则拒绝写入返回400。这与检查点重放一节提到的乱序/重复样本问题相关部署时需一并考虑。字段值类型契约源码级key_value指标数值在 PrometheusSerializer.createValueExtractor() 中支持STRING、INT、FLOAT、DOUBLE其中DOUBLE直接取值其余解析为数值字段为空时返回Double.NaN。文档明确推荐使用double类型字段可以规避解析与精度问题。key_label标签在 createLabelExtractor() 中要求字段 SQL 类型必须为MAP否则抛出UNSUPPORTED_DATA_TYPE异常值为null时返回空 Map。标签 Map 中建议包含__name__键用于表示指标名。key_timestamp时间戳支持四种类型转换规则如下表对应 createTimestampExtractor()字段类型处理方式timestamp按本地时区ZoneId.systemDefault()转换为毫秒级时间戳bigint直接按毫秒级时间戳处理double按 Unix 秒级时间戳处理乘以 1000 转为毫秒string按毫秒级时间戳字符串解析Long.parseLong未配置key_timestamp或字段值为null时统一回退到System.currentTimeMillis()即使用当前系统时间。接收器选项全解名称类型是否必填默认值描述urlString是-Prometheus 兼容 remote write API 地址如http://prometheus:9090/api/v1/writekey_labelString是-上游数据中保存 Prometheus 标签的字段名建议为 map 类型key_valueString是-上游数据中保存指标值的字段名推荐double类型key_timestampString否-时间戳字段名不配置时使用当前系统时间headersMap否-HTTP 请求头retryInt否-HTTP 请求出现IOException时的最大重试次数retry_backoff_multiplier_msInt否100重试退避时间倍数单位毫秒retry_backoff_max_msInt否10000最大重试退避时间单位毫秒batch_sizeInt否1024写入 Prometheus 前缓存的行数必须大于 0multi_table_sink_replicaInt否1多表写入时每张表使用的写入器副本数common-optionsConfig否-接收器插件通用参数见接收器通用选项选项的源码出处url、retry、retry_backoff_multiplier_ms、retry_backoff_max_ms、headers继承自 HTTP 基础模块 HttpCommonOptions.java其中url无默认值必填retry无默认值不配置即不启用重试退避倍数与上限默认分别为 100ms 与 10000mskey_label、key_value、key_timestamp、batch_size定义于 PrometheusSinkOptions.java其中batch_size默认1024配置装配逻辑见 PrometheusSinkConfig.loadConfig()。key_label指标标签与__name__key_label对应字段建议为mapstring, string会被直接转换为 Prometheus 标签label。强烈建议在 map 中包含__name__键作为指标名否则写入后序列series没有可观测的指标名。同时Sink 会自动补充 remote write 协议必需的三个请求头见 PrometheusSink 构造器Content-type: application/x-protobufContent-Encoding: snappyX-Prometheus-Remote-Write-Version: 0.1.0即使你通过headers自定义了请求头上述三个协议头也会被强制写入put覆盖保证请求可被后端正确解析。key_timestamp时间戳字段类型矩阵见上文字段值类型契约表格。一句话总结不配置就取当前系统时间配置后按timestamp/bigint/double/string四种类型自动归一化为毫秒级时间戳。multi_table_sink_replica多表写入并行度多表写入时每张表使用的 Sink Writer 副本数默认值为1。只有当单张表需要更高写入并行度时才建议调大多表场景下默认即可。这与SupportMultiTableSink的实现对应PrometheusSink.java。刷新机制三种触发时机与引擎级定时刷新Prometheus Sink 采用缓存批量发送模型write()先把序列化后的 Point 追加进batchList只有满足触发条件才真正发出 HTTP 请求。触发时机共有三种batch 触发缓存行数达到batch_size 0时生效见 PrometheusWriter.write()定时刷新Timer Flush由引擎按sink.flush.interval定时发出刷新信号检查点刷新prepareCommit()中无条件刷新关闭刷新close()中把剩余缓存全部发出。定时刷新仅 Zeta 引擎支持即使上游数据空闲、缓存行数还没达到batch_size接收器也可以按定时器刷新缓存把已缓存的采样点发送出去。该定时器由引擎驱动而不是由连接器驱动目前仅 SeaTunnel Zeta 支持。在作业的env中设置sink.flush.interval单位毫秒即可启用env { sink.flush.interval 10000 }实现层面PrometheusWriter 构造器 通过context.registerFlushAction(this::flush)向引擎注册刷新动作。引擎会在正常的 Sink 数据处理线程上触发刷新因此连接器不需要自己维护后台线程不会与写入、检查点、关闭等流程产生并发刷新失败会被抛给引擎而不会被静默丢弃。对应测试 PrometheusWriterTest.shouldRegisterFlushActionAndFlushBufferedRecordsOnSignal() 验证了注册 flush action 后单条写入只进缓存不发送只有引擎触发 signal 才真正doPost。Spark / Flink 上的行为差异Spark 和 Flink 的 Sink 写入器上下文并未实现registerFlushAction保留接口的 no-op 默认实现因此在它们之上没有检查点之间的定时刷新缓存会在以下时机被刷新达到batch_size检查点时prepareCommit()内刷新写入器关闭时close()内刷新。因此在这两个引擎上缓存的采样点最多保留一个检查点间隔而不会一直攒到batch_size或作业结束。如需降低 Spark/Flink 上检查点之间的延迟请相应调小batch_size。相关测试见 shouldFlushOnCloseWhenEngineNeverInvokesFlushAction() 与 shouldFlushOnPrepareCommitWhenEngineNeverInvokesFlushAction()。Zeta 上的双重触发检查点刷新在所有引擎上都会执行包括 Zeta。因此在 Zeta 上缓存会同时被sink.flush.interval和每个检查点触发刷新如果检查点间隔短于sink.flush.interval刷新会比仅靠定时器时更频繁每批更小。这是预期行为如果关注请求频率请同时调整sink.flush.interval和检查点间隔。旧版flush_interval选项的迁移提示连接器曾在自身配置中支持flush_interval该选项已被移除改为引擎级sink.flush.interval。升级后的作业配置若残留旧键PrometheusWriter 会打印警告日志提示改用env.sink.flush.interval仅 Zeta 生效Spark/Flink 无定时刷新请调batch_size避免运维人员以为周期刷新仍生效而实际被静默忽略。检查点刷新与失败处理检查点刷新是一次 remote-write 请求。prepareCommit()中调用flush()而 flush() 实现 只有在收到204 No Content时才清空缓存并返回其他任何响应码如400、5xx都会抛出PrometheusConnectorException错误码FLUSH_DATA_FAILED。因此瞬时失败会导致检查点失败。网络抖动、接收端重启或5xx响应都会让当前检查点失败。Flink 的tolerableCheckpointFailureNumber默认是0一次失败就会重启作业在 Spark 和 Flink 上对于低吞吐作业你可能需要调高引擎的可容忍检查点失败次数。测试 shouldPropagateFlushFailure() 验证了失败必须向上传播而非静默吞掉。连接器内部有界重试带退避的改进在社区 issue #11911 中跟踪当前flush()不做内部重试。重放的安全性取决于接收端。检查点失败后作业会重启Source 从上一次成功的检查点重放因此缓存的采样点会被重新发送。只有当 remote-write 接收端接受完全相同的重复样本相同的 labels、timestamp 和 value时这才是安全的如果接收端拒绝相同 timestamp 但 value 不同的样本或拒绝乱序样本Prometheus TSDB以及 Cortex、Mimir、Thanos 等接收端对这些情况返回400那么重放的刷新就会失败并持续让检查点失败。如果这对你的部署很重要请启用接收端的乱序窗口或确保重放的是完全相同的重复样本。整体交付语义是至少一次at-least-once而不是精准一次PrometheusWriter.prepareCommit() 注释。close()阶段同样保证最终刷新异常与 HTTP 客户端关闭异常同时发生时以刷新失败为主异常、其余作为 suppressed 附加避免真正的问题被掩盖close() 及测试 closeShouldKeepFlushExceptionWhenHttpClientCloseAlsoThrows()。HTTP 重试机制retry、retry_backoff_multiplier_ms、retry_backoff_max_ms三个参数控制 HTTP 请求失败时的重试行为实现在 HTTP 基础客户端 HttpClientProvider.java仅当retry 1时才构建 Retryer重试条件限定为IOExceptionretryIfException判断异常链中是否含 IOException停止策略为stopAfterAttempt(retry)即最多重试指定次数等待策略为指数退避由 multiplier 与 max 控制默认 100ms 起步、封顶 10000ms。由于 remote-write 请求体是幂等的重复样本时后端可安全重放合理设置retry可以缓解瞬时网络故障。安装依赖使用 Prometheus 连接器时需要安装connector-prometheus依赖可通过两种方式使用 SeaTunnel 发行包中的install-plugin.sh脚本安装从 Maven 中央仓库获取 artifactorg.apache.seatunnel:connector-prometheusgroup id 为org.apache.seatunnel。安装完成后将该连接器的 jar 放入$SEATUNNEL_HOME/connectors/目录或由插件发现机制加载即可在作业配置中以Prometheus为插件名使用。关于插件安装与依赖隔离的更多细节可参考连接器隔离依赖说明。完整示例FakeSource 写入 Prometheus以下示例来自 Prometheus.md使用内置 FakeSource 生成两个带标签的采样点经 Prometheus Sink 写入本地 Prometheusenv { parallelism 1 job.mode BATCH } source { FakeSource { schema { fields { c_map mapstring, string c_double double c_timestamp timestamp } } plugin_output fake rows [ { kind INSERT fields [{__name__ : metric_1}, 1.23, CURRENT_TIMESTAMP] }, { kind INSERT fields [{__name__ : metric_2}, 1.23, CURRENT_TIMESTAMP] } ] } } sink { Prometheus { plugin_input fake url http://prometheus:9090/api/v1/write key_label c_map key_value c_double key_timestamp c_timestamp batch_size 1 } }要点说明数据字段与采样点字段一一对应c_map→key_label含__name__指标名、c_double→key_value、c_timestamp→key_timestampbatch_size 1表示每行都立即发送适合演示生产环境可保持默认1024或按吞吐调整运行前请确认url指向的 Prometheus 已开启 remote write--enable-featureremote-write-receiver并可达。Prometheus 兼容 Remote Write 示例VictoriaMetrics若目标后端是 VictoriaMetrics 等兼容服务只需替换urlVictoriaMetrics 同样暴露/api/v1/write其余配置不变sink { Prometheus { plugin_input fake url http://victoria-metrics:8428/api/v1/write key_label c_map key_value c_double key_timestamp c_timestamp batch_size 5 } }变更日志连接器的完整版本变更记录请参见 connector-prometheus 变更日志包括选项调整如flush_interval移除、定时刷新引入等历史变更。实践要点总结字段映射key_label用mapstring, string且建议带__name__key_value用doublekey_timestamp可选不配则取系统当前时间。协议头Content-type: application/x-protobuf、Content-Encoding: snappy、X-Prometheus-Remote-Write-Version: 0.1.0由连接器自动补充无需手动设置。刷新时机batch_size默认 1024触发批量发送Zeta 上可用env.sink.flush.interval启用引擎级定时刷新所有引擎的检查点与关闭都会刷新剩余缓存。失败语义刷新失败会使检查点失败并重放缓存数据交付保证为至少一次若接收端拒绝乱序/重复样本需开启乱序窗口或保证重放样本完全一致。重试配置retry仅对IOException生效配合指数退避默认 100ms 起、10000ms 封顶缓解瞬时故障。【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价