Kafka 进、Kafka 出的实时 ETL 管道:用 Pathway 将多时区日期时间统一为时间戳 Kafka 进、Kafka 出的实时 ETL 管道用 Pathway 将多时区日期时间统一为时间戳【免费下载链接】pathwayPython ETL framework for stream processing, real-time analytics, LLM pipelines, and RAG.项目地址: https://gitcode.com/GitHub_Trending/pa/pathway本文导读本文基于仓库 [examples/projects/kafka-ETL](https://link.gitcode.com/i/ff21b1822ee13ee40b802b7905a78201) 示例项目完整讲解如何用 Pathway Live Data Framework 搭建一条 “Kafka in → 处理 → Kafka out” 的实时 ETL 管道两条数据源分别以不同时区广播日期时间字符串Pathway 从 Kafka 主题中抽取Extract数据、把带时区的日期时间字符串转换为毫秒级时间戳Transform再写回第三个 Kafka 主题Load。读完本文你将掌握 Pathway 的 pw.io.kafka.read / pw.io.kafka.write 用法、rdkafka 连接参数语义、日期时间解析与 concat_reindex 合并技巧以及如何用 docker compose 一键跑通整套流式示例并落地结果 CSV。项目目标与整体架构kafka-ETL是一个麻雀虽小、五脏俱全的流式 ETL 参考实现。它演示了最典型的实时数据管道形态流式平台Kafka作为数据总线Pathway 作为中间的无状态/有状态计算引擎输入输出都是 Kafka 主题数据无需落盘即可持续流转。根据 README示例共分三层流式平台由 Kafka 与 ZooKeeper 组成负责承载全部消息主题数据生产者Python 容器模拟两个位于不同时区的数据源不断向 Kafka 广播带时区的日期时间字符串ETL 处理端Python 容器使用 Pathway Live Data Framework 连接这些主题完成 E抽取→ T将日期时间转为时间戳→ L写入第三个主题全流程。全部容器通过docker-compose编排管理。数据流方向如下stream-producer ──(date: 字符串 message)──▶ topic: timezone1 ─┐ └──(date: 字符串 message)──▶ topic: timezone2 ─┴─▶ Pathway(ETL) ──▶ topic: unified_timestamps三个业务主题分别是timezone1、timezone2输入与unified_timestamps输出。Pathway 端的转换逻辑站在业务角度可以概括为把不同时区格式的2026-09-07 12:34:56.123456 -0400这类字符串统一换算成UTC 毫秒级时间戳供下游做时间对齐、窗口聚合或时序分析。为什么在 Pathway 中做 ETL示例的关键卖点在于Pathway 的输入输出都是实时流式连接器connector。从 kafka 连接器 Python 封装 的文档字符串可以看到pw.io.kafka.read的默认模式为streaming即引擎会持续等待新消息、随到随算并送入计算图这意味着同一个 Python 代码片段既完成了连接又完成了持续处理不再需要单独维护一套采集任务和一套批处理任务。下面的小节将结合真实源码逐一拆解。一键启动docker-compose 与 Makefile启动方式非常简单在 examples/projects/kafka-ETL 目录下执行# 方式一直接用 docker compose docker compose up -d # 方式二通过 Makefile make两种方式等价make内部执行的正是docker compose up -d见 Makefile。docker-compose 服务拓扑examples/projects/kafka-ETL/docker-compose.yml 共声明 4 个服务全部挂载在同一个etl-kafka-networkbridge 网络中服务镜像 / 构建来源职责依赖zookeeperconfluentinc/cp-zookeeper:5.5.3Kafka 依赖的协调服务ZOOKEEPER_CLIENT_PORT: 2181—kafkaconfluentinc/cp-enterprise-kafka:5.5.3消息代理承载全部主题zookeeperpathway本地构建pathway-src/DockerfileETL 主程序python etl.pykafkastream-producer本地构建producer-src/Dockerfile数据生产者python create-stream.pykafka两个自研容器的 Dockerfile 都以python:3.10为基础镜像pathway-src/Dockerfile 执行pip install -U pathway安装最新版 Pathway并把etl.py、read-results.py复制进容器容器入口为python etl.pyproducer-src/Dockerfile 安装kafka-python客户端入口为python create-stream.py。值得注意的 Kafka 服务环境变量见 docker-compose.ymlKAFKA_AUTO_CREATE_TOPICS: true允许主题自动创建因此timezone1、timezone2、unified_timestamps无需手工建主题KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9092容器间通过 compose 服务名kafka而非 IP 互通这正是各 Python 程序里bootstrap.servers: kafka:9092能生效的前提KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1单节点 Kafka 的必要配置默认 3 会因副本不足而失败。提示depends_on只保证容器先启动不保证 Kafka 服务已就绪。因此 etl.py 在pw.run()前显式time.sleep(20)create-stream.py 在真正发消息前也先time.sleep(30)用时间窗规避连接竞态。这正是单机演示型编排与生产环境应使用健康检查 / 重试机制的典型差别。第一站模拟双时区数据源Extract 侧的上游生产者容器由 producer-src/create-stream.py 驱动它使用标准库zoneinfo在America/New_York纽约与Europe/Paris巴黎两个时区下生成当前时间import json import random import time from datetime import datetime from zoneinfo import ZoneInfo from kafka import KafkaProducer input_size 100 random.seed(0) topic1 timezone1 topic2 timezone2 timezone1 ZoneInfo(America/New_York) timezone2 ZoneInfo(Europe/Paris) str_repr %Y-%m-%d %H:%M:%S.%f %z api_version (0, 10, 2) def generate_stream(): time.sleep(30) producer1 KafkaProducer( bootstrap_servers[kafka:9092], security_protocolPLAINTEXT, api_versionapi_version, ) producer2 KafkaProducer( bootstrap_servers[kafka:9092], security_protocolPLAINTEXT, api_versionapi_version, ) def send_message(timezone: ZoneInfo, producer: KafkaProducer, i: int): timestamp datetime.now(timezone) message_json {date: timestamp.strftime(str_repr), message: str(i)} producer.send(topic1, (json.dumps(message_json)).encode(utf-8)) for i in range(input_size): if random.choice([True, False]): send_message(timezone1, producer1, i) else: send_message(timezone2, producer2, i) time.sleep(1) time.sleep(2) producer1.close() producer2.close() if __name__ __main__: generate_stream()它的行为要点固定random.seed(0)保证可复现生成input_size 100条记录间隔 1 秒发送模拟每秒一条的持续数据流每条 JSON 消息只含两个字段date按%Y-%m-%d %H:%M:%S.%f %z格式化的字符串末尾%z携带时区偏移如-0400或0200和message消息序号字符串字符串格式str_repr在生产者与 ETL 端完全一致见 create-stream.py 与 etl.py这是上下游解析成功的前提。对双主题拓扑的一点源码观察按 README 的设计意图两条数据源应分别广播到timezone1与timezone2两个主题。不过从当前提交的 create-stream.py 源码看send_message内部固定将消息发往topic1topic2timezone2并未实际承接消息——随机分支的真正差异体现在date字段内嵌的时区偏移纽约-04xx与巴黎02xx以及使用的 producer 连接上。因此在观察本示例时重点应放在同一管道如何处理携带不同 UTC 偏移的日期字符串上若希望严格复现双主题分流只需把对应分支的producer.send(...)目标改为topic2即可Pathway 端的双read结构保持不变。第二站Pathway ETL 处理端核心代码逐行拆解ETL 的全部逻辑集中在单个文件 pathway-src/etl.py下面逐段解析。连接参数与 Schema 声明pw.set_license_key(demo-license-key-with-telemetry) rdkafka_settings { bootstrap.servers: kafka:9092, security.protocol: plaintext, group.id: 0, session.timeout.ms: 6000, auto.offset.reset: earliest, } str_repr %Y-%m-%d %H:%M:%S.%f %z class InputStreamSchema(pw.Schema): date: str message: str几个关键点pw.set_license_key(...)示例默认使用带遥测的 demo key 以启用高级特性使用 Pathway Community社区版时把这一行注释掉即可两处代码注释都做了同样说明。rdkafka_settings是 librdkafka 风格配置字典。pw.io.kafka.read/write的第一个参数要求Connection settings in the format of librdkafka见 python/pathway/io/kafka/init.py 的 API 文档这些键会被透传给底层的 rdkafkaRust 侧通过rdkafkacrate 封装见 src/connectors/data_storage/kafka.rs。各键含义如下配置键示例值作用bootstrap.serverskafka:9092初始 broker 地址列表用 compose 服务名在容器网络中解析security.protocolplaintext明文传输生产环境通常替换为SASL_SSL/SSLgroup.id0消费者组 ID。示例输入侧固定为0结果读取脚本则每次用uuid4()生成新组避免消费位点冲突session.timeout.ms6000会话超时broker 据此判定消费者是否失联auto.offset.resetearliest新消费者组无已提交位点时从最早消息开始读保证 ETL 不丢数据InputStreamSchema继承自pw.Schema声明 JSON 负载的两个字段date: str、message: str。Pathway 的 JSON 解析器会按 schema 中定义的列名从 JSON 对象取值生成表列。两个输入连接器pw.io.kafka.readtimestamps_timezone_1 pw.io.kafka.read( rdkafka_settings, topictimezone1, formatjson, schemaInputStreamSchema, autocommit_duration_ms100, ) timestamps_timezone_2 pw.io.kafka.read( rdkafka_settings, topictimezone2, formatjson, schemaInputStreamSchema, autocommit_duration_ms100, )两个read分别订阅timezone1与timezone2共用同一份InputStreamSchema。从 kafka 连接器源码 可确认此处用到的参数语义formatjson连接器先将消息负载按 JSON 解析再依据schema定义创建列另有raw与plaintext两种格式分别保留原始字节 / UTF-8 文本此时表包含key、data两列autocommit_duration_ms100相邻两次提交之间允许的最大时间间隔每隔这么多毫秒连接器会把收到的更新提交并送入 Pathway 计算图。其默认值为 1500ms见 python/pathway/io/kafka/init.py示例把它压低到100是为了让毫秒级延迟的小演示获得更平滑的实时感read默认工作于modestreaming持续等待新消息这正是Kafka out 结果会随新数据不断追加的底层原因。变换字符串日期时间 → 毫秒时间戳def convert_to_timestamp(table: pw.Table) - pw.Table: table table.select( datepw.this.date.dt.strptime(fmtstr_repr, contains_timezoneTrue), messagepw.this.message, ) table_timestamp table.select( timestamppw.this.date.dt.timestamp(unitms), messagepw.this.message, ) return table_timestamp两步变换分别对应 Pathway 日期时间表达式 API 中的两个方法实现见 python/pathway/internals/expressions/date_time.pydt.strptime(fmt..., contains_timezoneTrue)date_time.py把date字符串按%Y-%m-%d %H:%M:%S.%f %z解析为带时区语义的日期时间列。contains_timezoneTrue告知解析器格式中含时区偏移%z从而能正确识别-0400/0200这类偏移是跨时区统一的第一步dt.timestamp(unitms)date_time.py将解析后的本地时间含偏移换算为UTC 纪元毫秒时间戳。不同时区的同一物理时刻在这里被统一到同一条绝对时间轴上这是 TTransform阶段的核心产物。由于函数是Table → Table的纯列变换可对两个输入表复用timestamps_timezone_1 convert_to_timestamp(timestamps_timezone_1) timestamps_timezone_2 convert_to_timestamp(timestamps_timezone_2) timestamps_unified timestamps_timezone_1.concat_reindex(timestamps_timezone_2)concat_reindex定义见 python/pathway/internals/table.py将两张结构相同的表纵向拼接并重新生成行主键reindex即重排/再造索引从而规避两表原有主键冲突。此时两条不同时区的流已被合为一张统一时间戳 消息内容的宽表。输出连接器pw.io.kafka.writepw.io.kafka.write( timestamps_unified, rdkafka_settings, topic_nameunified_timestamps, formatjson ) # We wait for Kafka to be ready. time.sleep(20) # We launch the computation. pw.run()pw.io.kafka.write把处理结果以 JSON 格式持续写入unified_timestamps主题。需要明确上面只是声明了数据流图的读写两端真正的计算由最后一行pw.run()触发并阻塞运行持续监听、持续处理、持续写出time.sleep(20)与生产者侧的time.sleep(30)一样是等待 Kafka 就绪的防御式写法。至此完整的 ETL 调用链为pw.io.kafka.read(timezone1) ─┐ ├─ convert_to_timestamp ── concat_reindex ── pw.io.kafka.write(unified_timestamps) pw.io.kafka.read(timezone2) ─┘第三站读取结果并落地为 CSVREADME 提供了读取结果的完整操作路径。ETL 正常运行后unified_timestamps主题中已持续产出转换好的时间戳可通过 Pathway 提供的 pathway-src/read-results.py 把它读出来写入 CSV# 1) 进入运行 Pathway ETL 的容器 make connect # 等价于docker compose exec -it pathway bash # 2) 在容器内执行结果读取脚本 python read-results.py脚本本身仍然是一段迷你 Pathway 程序from uuid import uuid4 import pathway as pw pw.set_license_key(demo-license-key-with-telemetry) rdkafka_settings { bootstrap.servers: kafka:9092, security.protocol: plaintext, group.id: str(uuid4()), session.timeout.ms: 6000, auto.offset.reset: earliest, } topic_name unified_timestamps class InputStreamSchema(pw.Schema): timestamp: float message: str def read_results(): table pw.io.kafka.read( rdkafka_settings, topictopic_name, schemaInputStreamSchema, formatjson, autocommit_duration_ms100, ) pw.io.csv.write(table, ./results.csv) pw.run() if __name__ __main__: read_results()值得关注的两处设计细节读取侧 Schema 与 ETL 输出侧一一对应timestamp: float毫秒时间戳、message: str因为 ETL 写出的正是这两个 JSON 字段group.id str(uuid4())每次运行都使用全新消费者组配合auto.offset.reset: earliest保证即使之前消费过本次也会从最早消息读全不会因已提交位点而读不到历史数据。pw.io.csv.write(table, ./results.csv)把每条到达的行追加写入results.csv因为连接器默认处于streaming模式新到达的数据会被自动追加无需重启脚本这也正是 README 中 New data will be automatically added 的机制来源。Makefile 常用运维命令速查Makefile 集中封装了整套生命周期管理下表整理了全部目标Makefile 目标等价命令用途make/make builddocker compose up -d构建并后台启动全部容器make stopdocker compose down -v --remove-orphans 删除两个自研镜像停止并清理含卷与孤儿容器make connectdocker compose exec -it pathway bash进入 ETL 容器用于跑read-results.pymake connect-proddocker compose exec -it stream-producer bash进入生产者容器make connect-kafkadocker compose exec -it kafka bash进入 Kafka 容器可手工查主题make logsdocker compose logs pathway查看 ETL 日志make logs-proddocker compose logs stream-producer查看生产者日志make logs-kafkadocker compose logs kafka查看 Kafka 日志make logs-zookeeperdocker compose logs zookeeper查看 ZooKeeper 日志生产端每 1 秒产出一条消息、总计 100 条因此可在约 100 秒内观察results.csv被持续追加若想换用其它 Kafka 客户端直读unified_timestamps主题验证结果则不受限于示例脚本按 README 所说用你习惯的任何 Kafka 访问方式即可。底层支撑Kafka 连接器的实现侧面为便于读者继续深入这里补充几条源码线索Python 层 APIpython/pathway/io/kafka/init.py 定义了read/write的全部参数与缺省值如mode、format、autocommit_duration_ms、json_field_paths、with_metadata等还支持通过schema_registry_settings对接 Confluent Schema RegistryRust 层实现src/connectors/data_storage/kafka.rs 是 Kafka 读写器的核心实现基于rdkafkacrate 构建消费/生产实例并在启动阶段执行主题元数据获取与水位线watermark检查——例如连接器在启动时会校验主题是否存在、读取各分区 offset 低/高水位相关错误类型KafkaReaderErrorTopicNotFound、MetadataFetch 等会作为启动期错误抛回 Python 层见 kafka.rs集成测试仓库在 integration_tests/kafka 目录下提供了面向 Kafka 连接器的端到端集成测试可作为自行验证连接参数与读写行为时的参考。这种 Python API → 底层 Rust 连接器 的分层使得rdkafka_settings里的每个键都有明确的下游消费方配置错误如 broker 不可达、主题不存在通常会在pw.run()启动阶段以清晰的错误信息暴露出来便于快速定位。小结从示例到生产的最小迁移清单本示例用约 60 行核心代码证明了以 Kafka 为主题的流式 ETL在 Pathway 中的简洁形态read声明输入、若干列变换函数描述 T、write声明输出、pw.run()启动常驻计算。在动手把该示例迁移到真实场景时建议对照以下清单将security.protocol从plaintext升级为与 broker 匹配的SASL_SSL等并补充认证凭据把固定group.id改为按消费者语义规划的组名避免多副本 ETL 实例间位点互相干扰视时延/吞吐权衡调整autocommit_duration_ms默认 1500ms 已适合多数场景用 broker 健康检查替代示例中的time.sleep等待并让 ETL 具备断线重连的容错能力若消息量较大可通过parallel_readers、with_metadata等参数见 kafka 连接器 API做并行消费与元数据追踪。掌握以上要点后你便可以在自己的 Kafka 生态中复用这套 Kafka in / Pathway 处理 / Kafka out 的管道骨架把本文的跨时区时间戳统一化替换为任意 ETL 变换逻辑。【免费下载链接】pathwayPython ETL framework for stream processing, real-time analytics, LLM pipelines, and RAG.项目地址: https://gitcode.com/GitHub_Trending/pa/pathway创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考