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,其余消息被静默丢弃。
本配方对应的模式定义为:
- Pattern:Kafka Patterns - Content-Based Routing
- 难度:Basic(基础)
- 核心组件:
kafka_franz(输入/输出)、mapping(Bloblang 映射) - 典型用例:根据消息内容字段将 Kafka 消息路由到不同 Topic
这种模式的关键价值在于:在不破坏消息顺序与分区语义的前提下完成过滤分流。如果只是简单地把消息读出来过滤再写入,很容易丢失源消息的分区键、时间戳或 Header,导致下游出现顺序错乱、分区分布变化等问题。本配方通过元数据保留 + 手动分区机制,精确规避了这些坑。
完整配置与逐段解析
完整配置位于仓库内的 content-based-router.yaml,整个管线分为三大部分:输入(kafka_franz+ 两个 processors)、输出(kafka_franz),全部字段均以环境变量注入,不硬编码任何凭据。
输入配置:消费源 Topic
input: 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_topics | false | 是否将topics解释为正则表达式,关闭即按字面 Topic 名匹配 |
consumer_group | ${CONSUMER_GROUP} | 消费组;指定后分区会在同一消费组的客户端之间自动均衡 |
auto_replay_nacks | true | 消息处理失败(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: "manual" | manual | 显式指定分区策略,配合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_write | true | 开启幂等写入,降低重复投递 |
max_message_bytes | 1024 | 压缩前的批大小(该值可按实际负载调整) |
broker_write_max_bytes | 100MiB | 单次请求的最大字节数,支撑大消息 |
max_in_flight | 256 | 并行写入的批次数量上限,提升吞吐 |
client_id | content_based_router | 客户端标识,便于追踪与调试 |
其中partitioner: "manual"的语义可以从源码得到精确佐证。在 franz_writer.go 中,kafka_franz输出支持的四种分区器分别为:
murmur2_hash:默认的 murmur2 哈希分区(kgo.StickyKeyPartitioner);round_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: "manual"+partition插值直接把源分区号作为目标分区号,比哈希更直接——只要源 Topic 与目标 Topic 的分区数一致,消息会精确落在"同号分区"。 - 时间戳还原:
timestamp恢复原始事件时间,避免消费-重写过程把事件时间替换为处理时间,这对下游的时间窗口聚合、延迟统计等至关重要。
除此之外,auto_replay_nacks: true保证失败消息会被重放重试而非直接丢弃,配合idempotent_write: true尽量消除重试带来的重复投递,进一步稳定了顺序与去重语义。
本地测试与验证
设置环境变量并运行管线
# Set environment variables export KAFKA_BROKER=localhost:9092 export SOURCE_TOPIC=test_in export DEST_TOPIC=topic_a export CONSUMER_GROUP=test_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.yaml>lint会校验 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/else+deleted()按字段值丢弃不匹配消息; - 后还原:输出端以
partitioner: "manual"+ 插值表达式还原分区号、消息键与时间戳,并用include_patterns: [".*"]保留全部自定义 Header。
这套"备份—过滤—还原"的组合,配合幂等写入与自动重放,既完成了内容驱动的消息分流,又从根源上保住了分布式系统中珍贵的顺序保证,是生产级 Kafka 消息路由的可靠范式。
【免费下载链接】connectFancy stream processing made operationally mundane项目地址: https://gitcode.com/GitHub_Trending/con/connect
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考