Redpanda Connect 基于内容的 Kafka 消息路由器(Content-Based Router)完整实战指南
Redpanda Connect 基于内容的 Kafka 消息路由器Content-Based Router完整实战指南【免费下载链接】connectFancy stream processing made operationally mundane项目地址: https://gitcode.com/GitHub_Trending/con/connect导读本文讲解 Redpanda Connect本仓库项目中Content-Based Router基于内容的路由这一经典 Kafka 集成模式如何通过 Bloblang 检查消息负载字段、只将满足条件的消息转发到目标 Topic同时完整保留分区键、分区号、时间戳与 Header从而维持分布式系统中的消息顺序与分区语义。读完本文你将掌握完整可运行的配置文件、手动分区manual partitioner的底层原理、元数据保留的关键实现以及多目标路由switch扩展方案可直接用于生产级 Kafka 消息过滤与分流场景。本文主体素材来自仓库中的官方配方文档 content-based-router.md 与其配套的完整配置 content-based-router.yaml两者位于pipeline-assistant技能的生产配方目录是经过校验的可用配置模板。模式概览什么是基于内容的路由Content-Based Router基于内容的路由器是一种消息路由模式系统根据消息内容负载字段动态决定把消息送往哪个目的地。在 Kafka 场景下最常见的形态是按字段过滤 单目标转发——即从源 Topic 消费消息检查负载中的某个字段例如marketid只有匹配特定值的消息被写入目标 Topic其余消息被静默丢弃。本配方对应的模式定义为PatternKafka Patterns - Content-Based Routing难度Basic基础核心组件kafka_franz输入/输出、mappingBloblang 映射典型用例根据消息内容字段将 Kafka 消息路由到不同 Topic这种模式的关键价值在于在不破坏消息顺序与分区语义的前提下完成过滤分流。如果只是简单地把消息读出来过滤再写入很容易丢失源消息的分区键、时间戳或 Header导致下游出现顺序错乱、分区分布变化等问题。本配方通过元数据保留 手动分区机制精确规避了这些坑。完整配置与逐段解析完整配置位于仓库内的 content-based-router.yaml整个管线分为三大部分输入kafka_franz 两个 processors、输出kafka_franz全部字段均以环境变量注入不硬编码任何凭据。输入配置消费源 Topicinput: label: consume_from_source kafka_franz: seed_brokers: [${KAFKA_BROKER}] topics: [${SOURCE_TOPIC}] regexp_topics: false consumer_group: ${CONSUMER_GROUP} auto_replay_nacks: true # Retry failed messages各字段说明字段值说明seed_brokers${KAFKA_BROKER}用于建立连接的 broker 地址列表支持逗号展开多个地址topics${SOURCE_TOPIC}要消费的源 Topic 列表regexp_topicsfalse是否将topics解释为正则表达式关闭即按字面 Topic 名匹配consumer_group${CONSUMER_GROUP}消费组指定后分区会在同一消费组的客户端之间自动均衡auto_replay_nackstrue消息处理失败nack时自动重放重试kafka_franz是使用 Franz Kafka 客户端库franz-go的 Kafka 输入组件完整字段文档见 inputs/kafka_franz.adoc。从该文档可以确认kafka_franz输入会给每条消息附加如下元数据字段kafka_key、kafka_topic、kafka_partition、kafka_offset、kafka_lag、kafka_timestamp_ms、kafka_timestamp_unix、kafka_tombstone_message以及全部记录头。这些元数据正是后续保留分区与顺序的关键原料。在输入节点之下紧跟两个 processors分别完成元数据保留和内容过滤。处理器一先备份 Kafka 元数据processors: # Preserve Kafka metadata before processing - label: copy_kafka_metadata mapping: | # Separate Kafka-specific metadata from custom metadata # This allows us to restore partition/key/timestamp in output let kafka_meta .filter(kv - kv.key.has_prefix(kafka_)) meta .filter(kv - !kv.key.has_prefix(kafka_)) meta kafka_metadata $kafka_meta这是整个配方中最精妙的一步逻辑拆解如下是 Bloblang 中访问全部元数据的引用。.filter(kv - kv.key.has_prefix(kafka_))把键名以kafka_开头的系统元数据如kafka_key、kafka_partition、kafka_timestamp_unix整体提取到一个局部变量$kafka_meta中。meta .filter(kv - !kv.key.has_prefix(kafka_))将剩余的自定义元数据业务 Header保留为常规元数据。meta kafka_metadata $kafka_meta把提取出的整份 Kafka 系统元数据重新打包到名为kafka_metadata的元数据键下作为 JSON 对象整体携带。为什么要先备份再还原原因在于输出端的kafka_franz组件对某些kafka_前缀元数据会有特殊处理例如把kafka_key之类字段作为消息键或分区依据使用在消息流转过程中直接依赖这些活元数据并不可靠。将其固化为独立元数据对象后在输出端可以稳定地通过${!metadata(kafka_metadata).kafka_partition}这类插值表达式还原分区键、分区号与时间戳。这也解释了为什么输出端的插值路径都写成metadata(kafka_metadata).xxx而不是直接使用顶层kafka_元数据。处理器二按字段内容过滤# Filter messages based on content - label: filter_by_marketid mapping: | # Route only NYSE messages if (this.marketid nyse) { root this } else { # Filter out non-NYSE messages root deleted() }过滤逻辑使用 Bloblang 映射this引用消息的 JSON 负载this.marketid读取marketid字段。当marketid nyse时root this保留整条消息原样通过否则执行root deleted()彻底删除丢弃该消息不会进入输出端。deleted()是 Bloblang 内置函数用于把当前消息标记为已删除被删除的消息不会继续传递也不会被写入输出。这样非匹配消息就被静默过滤掉了。从仓库源码可以印证这一点kafka_franz输入的元数据写入逻辑位于 franz_reader.go其中正是通过msg.MetaSetMut(kafka_key, ...)、msg.MetaSetMut(kafka_partition, ...)、msg.MetaSetMut(kafka_timestamp_unix, record.Timestamp.Unix())挂载这些 Kafka 元数据字段与本配方先备份、后还原的用法完全对应。输出配置手动分区 元数据还原output: label: write_to_destination kafka_franz: seed_brokers: [${KAFKA_BROKER}] topic: ${DEST_TOPIC} # Preserve source partition (maintains ordering) partitioner: manual partition: ${!metadata(\kafka_metadata\).kafka_partition} # Preserve source message key (maintains co-partitioning) key: ${!metadata(\kafka_metadata\).kafka_key} # Preserve source timestamp (maintains event time) timestamp: ${!metadata(\kafka_metadata\).kafka_timestamp_unix} # Preserve all custom headers metadata: include_patterns: [.*] # Use idempotent writes to minimize duplicates idempotent_write: true # Performance tuning max_message_bytes: 1024 # Batch size before compression broker_write_max_bytes: 100MiB # Max request size for large messages max_in_flight: 256 # High parallelism for throughput # Set client ID for tracing/debugging client_id: content_based_router关键配置项与作用字段值作用partitioner: manualmanual显式指定分区策略配合partition字段手动控制每条消息的落分区partition${!metadata(kafka_metadata).kafka_partition}从备份元数据还原源消息所在分区号key${!metadata(kafka_metadata).kafka_key}还原源消息的消息键维持基于键的共分区co-partitioning语义timestamp${!metadata(kafka_metadata).kafka_timestamp_unix}还原源消息的原始事件时间metadata.include_patterns[.*]将全部自定义 Header 作为消息头写入目标消息idempotent_writetrue开启幂等写入降低重复投递max_message_bytes1024压缩前的批大小该值可按实际负载调整broker_write_max_bytes100MiB单次请求的最大字节数支撑大消息max_in_flight256并行写入的批次数量上限提升吞吐client_idcontent_based_router客户端标识便于追踪与调试其中partitioner: manual的语义可以从源码得到精确佐证。在 franz_writer.go 中kafka_franz输出支持的四种分区器分别为murmur2_hash默认的 murmur2 哈希分区kgo.StickyKeyPartitionerround_robin轮询分区least_backup写入备份最少的分区manual手动选择分区要求同时配置partition字段对应kgo.ManualPartitioner。当partitioner设置为manual时每条消息通过${!metadata(...)...}插值表达式解析出整数分区号消息被精确写入与源消息相同的分区。由于同一分区内的消息天然保持写入顺序这样原分区号 原消息键的组合就从根源上维持了源 Topic 到目标 Topic 的顺序保证。输出端partition字段的官方语义同样可以在 outputs/kafka_franz.adoc 中查到该字段仅在partitioner为manual时生效插值结果必须是合法整数示例即${! meta(partition) }。此外需要留意一个与幂等写入相关的约束仓库源码 franz_writer.go 明确校验了idempotent_write开启时acks必须为all且max_in_flight_requests必须为1否则会直接报错。原因是幂等写入依赖每个 broker 单飞行请求来维持生产者序列号连续。这里的max_in_flight: 256与max_in_flight_requests是不同维度的参数前者表示并行写入的消息批数量后者表示单连接上的在途 produce 请求数。配方默认配置已符合该约束采用默认的acks: all与默认在途请求数 1使用时应避免将idempotent_write: true与较大的max_in_flight_requests同时设置否则管线会启动失败并持续重试。顺序保证为什么这套配置能保住消息顺序在分布式流处理中Kafka 只保证同一分区内消息的顺序。跨分区或跨 Topic 的顺序本身没有全局保证。因此路由后顺序不乱的正确含义是源 Topic 同一分区的消息写入目标 Topic 时仍落在同一分区且保持原有先后关系。本配方通过三个层次的协作实现这一目标消息键还原key插值恢复源消息键。在 Kafka 语义中同一键的消息被哈希到同一分区键的还原保证了基于键的分区归属不变。分区号还原partitioner: manualpartition插值直接把源分区号作为目标分区号比哈希更直接——只要源 Topic 与目标 Topic 的分区数一致消息会精确落在同号分区。时间戳还原timestamp恢复原始事件时间避免消费-重写过程把事件时间替换为处理时间这对下游的时间窗口聚合、延迟统计等至关重要。除此之外auto_replay_nacks: true保证失败消息会被重放重试而非直接丢弃配合idempotent_write: true尽量消除重试带来的重复投递进一步稳定了顺序与去重语义。本地测试与验证设置环境变量并运行管线# Set environment variables export KAFKA_BROKERlocalhost:9092 export SOURCE_TOPICtest_in export DEST_TOPICtopic_a export CONSUMER_GROUPtest_cg # Run the pipeline rpk connect run content-based-router.yaml生产测试消息并验证过滤# Produce test messages echo {marketid:nyse,symbol:AAPL,price:150} | rpk topic produce $SOURCE_TOPIC echo {marketid:nasdaq,symbol:MSFT,price:300} | rpk topic produce $SOURCE_TOPIC echo {marketid:nyse,symbol:GOOGL,price:2800} | rpk topic produce $SOURCE_TOPIC # Check output topic (only NYSE messages should appear) rpk topic consume $DEST_TOPIC预期结果三条消息中只有marketid为nyse的两条AAPL、GOOGL会出现在topic_a中MSFT 被静默丢弃。使用 lint 校验配置仓库的pipeline-assistant技能提供了一整套开发工作流见 SKILL.md其中与本配方最相关的是配置校验rpk connect lint [--env-file .env] pipeline.yamllint会校验 YAML 语法、组件配置与 Bloblang 表达式输出带具体位置的错误信息退出码 0 表示通过。配方目录下的 validate.sh 展示了批量校验脚本的写法它先加载.env.validation中的环境变量再对目录下每个*.yaml执行rpk connect lint任一文件校验失败即中止并打印错误——你可以在自己的项目中复用同样的模式把配置即代码纳入 CI。更稳妥的本地验证方式是用stdin/stdout先测试路由逻辑本身把输入替换为stdin、输出替换为stdout用echo管道直接运行快速确认 Bloblang 条件与deleted()行为是否符合预期再接入真实 Kafka。注意这种方式无法验证批处理、连接重试、顺序保证与并行处理等运行时行为生产部署前仍需真实集成测试。变体多目标路由switch 输出单目标过滤只是基于内容路由的基础形态。若需要把不同内容路由到不同 Topic可用switch输出替换过滤处理器output: switch: cases: - check: json(marketid) nyse output: kafka_franz: topic: topic_nyse - check: json(marketid) nasdaq output: kafka_franz: topic: topic_nasdaqswitch按顺序从上到下评估各分支的check条件这里用json(marketid)查询负载字段命中即写入对应 Topic。注意两点switch默认只执行第一个匹配的分支若希望一条消息同时落入多个分支需显式配置如switch.strict/多分支命中相关选项。分支输出各自需要独立的kafka_franz配置若要保持源分区与顺序同样要在各分支中应用partitioner: manual 元数据还原的写法。switch多目标路由是 CDC 复制等高级场景的基础仓库中的 cdc-replication.md 正是基于 switch 的进阶路由配方值得对照阅读。变体先经 stdin/stdout 快速验证路由逻辑pipeline-assistant技能文档还给出了一种无需 Kafka 的轻量验证法用stdin输入 元数据打标 switch输出到stdout在接入真实系统前验证路由决策。其核心思路是先用 mapping 处理器根据消息类型设置route元数据再用switch按meta(route)分发input: stdin: {} pipeline: processors: - mapping: | root this # Route based on message type if this.type error { meta route dlq } else if this.priority high { meta route urgent } else { meta route standard } output: switch: cases: - check: meta(route) dlq output: stdout: {} processors: - mapping: root DLQ: content().string() - check: meta(route) urgent output: stdout: {} processors: - mapping: root URGENT: content().string() - check: meta(route) standard output: stdout: {} processors: - mapping: root STANDARD: content().string()验证命令与预期输出echo {type:error,msg:failed} | rpk connect run test.yaml # Output: DLQ: {type:error,msg:failed} echo {priority:high,msg:urgent} | rpk connect run test.yaml # Output: URGENT: {priority:high,msg:urgent} echo {priority:low,msg:normal} | rpk connect run test.yaml # Output: STANDARD: {priority:low,msg:normal}这种先映射打标、后 switch 分发的写法与本配方的先映射过滤、后手动分区写入互为补充前者适合验证路由判定逻辑后者用于生产级 Kafka 落盘。重要细节与生产注意事项安全broker 地址等敏感信息一律使用环境变量${KAFKA_BROKER}、${SOURCE_TOPIC}、${DEST_TOPIC}、${CONSUMER_GROUP}注入配置文件本身不含任何明文凭据。如需更多密钥SASL 密码、API Key 等同样采用${VAR}语法并配合rpk connect run --env-file .env加载.env文件.env应加入.gitignore仅提交包含占位值的.env.example。性能max_in_flight: 256提供高并行度显著提升吞吐idempotent_write: true防止重试导致重复消息但如前文所述它要求acks: all且每 broker 在途请求数为 1这是吞吐与去重之间的权衡点broker_write_max_bytes: 100MiB允许单请求承载大消息适合负载较大的场景。错误处理auto_replay_nacks: true让失败消息进入重放重试而不是悄悄丢失。若需要更精细的失败处理如路由失败进入死信队列可参考 dlq-basic.md。顺序与共分区手动分区 键还原是保住顺序的根基若目标 Topic 分区数与源不一致跨分区的相对顺序在 Kafka 语义下本就不作保证设计时应确保目标 Topic 分区数 ≥ 源 Topic 分区数。相关配方与延伸阅读本配方属于pipeline-assistant技能生产配方库recipes 目录可与以下内容对照学习DLQ Basic死信队列处理路由失败的消息CDC Replication变更数据捕获复制基于 switch 的高级路由Multicast多播扇出一对多目标扇出kafka_franz 输入组件文档完整字段说明与元数据清单kafka_franz 输出组件文档partitioner、partition、metadata等字段语义。若要深入组件实现可继续阅读 franz_reader.go输入元数据挂载、franz_writer.go分区器、幂等写入与在途请求约束以及 input_kafka_franz.go输入组件注册与元数据文档。小结基于内容的路由是 Kafka 流处理中最常用、也最容易在顺序与分区语义上出错的模式之一。本配方给出的解法可以总结为三步先备份用 Bloblang 把kafka_前缀的系统元数据打包成kafka_metadata元数据对象再过滤用if/elsedeleted()按字段值丢弃不匹配消息后还原输出端以partitioner: manual 插值表达式还原分区号、消息键与时间戳并用include_patterns: [.*]保留全部自定义 Header。这套备份—过滤—还原的组合配合幂等写入与自动重放既完成了内容驱动的消息分流又从根源上保住了分布式系统中珍贵的顺序保证是生产级 Kafka 消息路由的可靠范式。【免费下载链接】connectFancy stream processing made operationally mundane项目地址: https://gitcode.com/GitHub_Trending/con/connect创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考