Pathway 实现 Kafka ETL:用 Python 统一跨时区事件流并写回 Kafka
本文以 Pathway(Pathway Live Data Framework,Python 流处理框架 + Rust 引擎)为技术主体,完整讲解如何用它构建一个 Kafka in、Kafka out 的 ETL 管道:通过 Kafka 输入连接器从多个 topic 抽取(Extract)带不同时区的日志、用 datetime 表达式统一转换(Transform)为时间戳、借助 concat_reindex 合并数据流,最后经 Kafka 输出连接器装载(Load)到一个合并 topic。读完本文,你将掌握 Pathway 连接器的配置方法、rdkafka_settings 的常用参数含义,以及时间处理与表拼接的底层实现,并可直接复刻一份可运行的完整工程。
业务场景:不同时区服务器日志的清洗任务
假设你入职了一家反欺诈(fraud-detection)公司:公司监控不同服务器的日志,一旦发现可疑模式就触发告警。你的职责是管理数据,保证交给数据科学团队的数据是干净、可直接使用的。
起初服务器都位于纽约,日志时间统一带有东部标准时间(EST)标记:
2024-02-05 10:01:52.884548 -0500
当公司把监控扩展到巴黎的服务器后,时间不再统一了——巴黎日志带的是中欧时间(CET)偏移:
2024-02-05 16:02:34.934749 +0100
同一个时点出现了 -0500 与 +0100 两种偏移,直接做时序对齐、窗口聚合或模型训练都会出错。为了维护数据一致性,必须把来自不同数据源的、时区各异的时间统一成一种与时区无关的标准格式。
这正是 ETL(Extract-Transform-Load)的典型用武之地:先 抽取(Extract) 数据,再 转换(Transform),最后 装载(Load)。本场景中:
- (E) 用 Pathway 的 Kafka 输入连接器从不同 topic 抽取多条数据流;
- (T) 用
datetime相关表达式把带时区的时间字符串解析并转换为时间戳(timestamp,与具体时区无关的绝对时刻); - (T) 用拼接函数
concat_reindex合并转换后的数据流; - (L) 用 Kafka 输出连接器把最终数据流装载回一个统一的 Kafka topic。
场景对应仓库中的完整工程位于 examples/projects/kafka-ETL,其中的 README.md 对项目组织结构做了简要说明。
ETL 架构:Kafka in,Kafka out
日志被发送到两个互不相同的 Kafka topic(每个时区一个)。Pathway 在此扮演 ETL 引擎:连接两个 topic、抽取数据、完成时区转换、把两条结果数据流拼接为一条,最后把结果写入第三个 Kafka topic(unified_timestamps)。整体数据流是"Kafka 进、Kafka 出"。
Docker 容器划分
工程通过 Docker 编排四个容器:
- Kafka:事件流平台本体;
- Zookeeper:Kafka 的协调服务;
- Pathway:承载 ETL 逻辑的容器;
- producer(stream-producer):模拟公司服务器、按秒产生日志数据并投递到 Kafka 的容器。
Kafka 与 Zookeeper 直接在 docker-compose.yml 中管理,Pathway 与 producer 则使用各自的 Dockerfile 构建。从仓库可见最终目录结构:
.
├── pathway-src/
│ ├── Dockerfile
│ ├── etl.py
│ └── read-results.py
├── producer-src/
│ ├── create-stream.py
│ └── Dockerfile
├── docker-compose.yml
└── Makefile
在 docker-compose.yml 中可以看到各服务的细节:
zookeeper使用confluentinc/cp-zookeeper:5.5.3,暴露2181端口;kafka使用confluentinc/cp-enterprise-kafka:5.5.3,声明KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9092、KAFKA_AUTO_CREATE_TOPICS: true(允许 topic 自动创建,示例因此无需手工建 topic)等环境变量,并把9092映射到宿主机;pathway与stream-producer分别通过./pathway-src/Dockerfile、./producer-src/Dockerfile构建,且都depends_on: [kafka]。
pathway-src/Dockerfile 很简单:基于 python:3.10,执行 pip install -U pathway 后把两个脚本拷入容器,默认命令为 python etl.py。
工程还提供了 Makefile 封装常用操作,例如:
make build:docker compose up -d拉起全部容器;make stop:docker compose down -v --remove-orphans并清理镜像;make connect/make connect-prod/make connect-kafka:分别进入 pathway、producer、kafka 容器;make logs/make logs-prod:查看 pathway / producer 容器日志。
数据生成:producer 按秒产生日志
数据由 Python 脚本生成:每隔 1 秒产生一条携带当前 datetime 的新日志,随机绑定两个时区之一并投递到对应的 Kafka topic;每条消息还带一个 message 字段记录日志编号,方便核对顺序。核心逻辑见 create-stream.py:
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) # 等 Kafka 就绪
producer1 = KafkaProducer(
bootstrap_servers=["kafka:9092"],
security_protocol="PLAINTEXT",
api_version=api_version,
)
producer2 = KafkaProducer(
bootstrap_servers=["kafka:9092"],
security_protocol="PLAINTEXT",
api_version=api_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()
要点说明:
- 两个
KafkaProducer分别对应两个时区,消息体中date字段用strftime以%Y-%m-%d %H:%M:%S.%f %z格式(含%z时区偏移)序列化,message字段记录编号; bootstrap_servers=["kafka:9092"]使用 compose 网络内的服务名kafka而非localhost;- 每条消息体是 JSON 字符串(
json.dumps后utf-8编码),这决定了后续 Pathway 侧要用format="json"解析; - 随机种子
random.seed(0)与input_size = 100使演示可复现、可收敛。
Extract:用 Pathway Kafka 输入连接器抽取数据
Pathway 提供了连接器(connector)来对接各类数据源。对于 Kafka,可以用 pw.io.kafka.read 建立输入连接。
在 Pathway 中,数据以**表(Table)的形式表示,抽取前需要先用模式(Schema)**声明数据的列与类型。这里每条日志包含两个字符串字段:
class InputStreamSchema(pw.Schema):
date: str
message: str
由于有两个输入 topic,需要每个 topic 一个连接器,但它们可以共享同一份连接设置:
rdkafka_settings = {
"bootstrap.servers": "kafka:9092",
"security.protocol": "plaintext",
"group.id": "0",
"session.timeout.ms": "6000",
"auto.offset.reset": "earliest",
}
timestamps_timezone_1 = pw.io.kafka.read(
rdkafka_settings,
topic="timezone1",
format="json",
schema=InputStreamSchema,
autocommit_duration_ms=100,
)
timestamps_timezone_2 = pw.io.kafka.read(
rdkafka_settings,
topic="timezone2",
format="json",
schema=InputStreamSchema,
autocommit_duration_ms=100,
)
rdkafka_settings 以 librdkafka 配置的格式给出,本场景中各键含义如下:
| 配置键 | 取值 | 作用 |
|---|---|---|
bootstrap.servers |
kafka:9092 |
Kafka 集群入口地址(compose 网络内用服务名) |
security.protocol |
plaintext |
明文传输,不启用 TLS/SASL |
group.id |
0 |
消费组 ID,决定消费位点的管理与隔离 |
session.timeout.ms |
6000 |
会话超时(毫秒),用于故障检测 |
auto.offset.reset |
earliest |
无已提交位点时从最早消息开始消费 |
可以对照 Kafka 输入连接器实现 理解各参数:read 的关键形参包括 topic(可传单个 topic 名或列表)、schema(结果表模式)、mode(默认 "streaming",持续等待新消息实时处理;设为 "static" 则只处理启动时已存在的数据)、format(支持 "raw" / "plaintext" / "json",其中 "raw" 与 "plaintext" 输出 key/data 两列,"json" 则会按 schema 解析 JSON 负载字段并生成对应列)、以及 autocommit_duration_ms(两次提交之间的最大间隔,每次提交会把增量更新推入计算图,默认 1500)。
实现中还校验了必须提供非空的 bootstrap.servers 与 group.id(见 python/pathway/io/kafka/init.py),并支持 start_from_timestamp_ms、parallel_readers、with_metadata、json_field_paths 等进阶参数,用于指定消费起点、并行读、读取消息元数据或按 JSON Pointer 重映射字段。
由于 Kafka 消息中 date 是形如 2024-02-05 10:01:52.884548 -0500 的字符串,这里 format="json" + schema 会直接把 date、message 两个 JSON 字段解析为表列。
Transform:统一时间并拼接为单一数据流
拿到日志后要做的转换核心是"时间":把带时区偏移的时间字符串变成与时钟时区无关的时间戳。Pathway 在表达式层提供了完整的 datetime 命名空间,见 python/pathway/internals/expressions/date_time.py。
转换函数如下:
def convert_to_timestamp(table):
table = table.select(
date=pw.this.date.dt.strptime(fmt=str_repr, contains_timezone=True),
message=pw.this.message,
)
table_timestamp = table.select(
timestamp=pw.this.date.dt.timestamp(unit="ms"),
message=pw.this.message,
)
return table_timestamp
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)
分三步理解这段转换:
第 1 步:strptime 把字符串解析成 datetime。 调用 pw.this.date.dt.strptime(fmt=str_repr, contains_timezone=True),fmt 对应消息里的格式化串 "%Y-%m-%d %H:%M:%S.%f %z";contains_timezone=True 告知解析器输入带时区偏移,%z 解析出的时区会自动被识别并标准化到统一的内部表示。方法签名可见 date_time.py 中的 strptime。
第 2 步:timestamp 把 datetime 换算成绝对时间戳。 pw.this.date.dt.timestamp(unit="ms") 输出以毫秒为单位、与偏移无关的时间戳(epoch 毫秒)。因为输入已经携带时区信息,无论原来写的是 -0500 还是 +0100,同一物理时刻都会得到同一个数值——时区差异在这一步被彻底抹平。实现见 date_time.py 中的 timestamp。
第 3 步:concat_reindex 拼接两条数据流。 两个时区分别转换成时间戳表后,用 timestamps_timezone_1.concat_reindex(timestamps_timezone_2) 合并为一张统一表。concat_reindex 会把多张同构表的行合并并对主键重新编号,避免来自不同流的行主键冲突。其定义位于 python/pathway/internals/table.py,属于表级拼接算子。
这是一个很朴素的示例。Pathway 还支持更复杂的有状态与时间态运算,例如基于 groupby/reduce 的分组聚合、基于窗口(windows)的时间窗口计算、ASOF join / interval join 这类时间对齐连接等,可作为后续扩展方向。
Load:把统一数据流写回 Kafka
转换完成后,还需要把结果送回 Kafka。使用 Kafka 输出连接器即可:
pw.io.kafka.write(
timestamps_unified, rdkafka_settings, topic_name="unified_timestamps", format="json"
)
由于写回的是同一个 Kafka 实例,rdkafka_settings 与输入连接器完全相同。pw.io.kafka.write 的实现见 python/pathway/io/kafka/init.py,它同样要求提供非空的 bootstrap.servers;topic_name 指定目标 topic(unified_timestamps),format="json" 决定以 JSON 形式序列化每行消息。
Run:让管道真正跑起来
注意:到目前为止,代码只是在构建管道——声明连接器与各种算子。此时系统里还没有实际数据流动;若现在就运行,管道虽已建立但不会发生任何计算。要开始摄取数据、驱动整条管道运行,必须调用 Pathway 的 run 函数:
pw.run()
调用之后输入连接器才会真正连接 Kafka 并持续装载数据,输出连接器也会把结果持续写入目标 topic。
Pathway 的计算核心是 Rust 引擎:数据经过连接器进入系统后,由 Rust 引擎多线程并发执行算子更新,Python 只是描述计算图的前端,因此不必受 GIL 等 Python 常规性能边界的限制,这也使整条 Kafka ETL 能以流式、低延迟的方式持续运行。仓库中 etl.py 的最终形态还包含两处工程化细节:
pw.set_license_key(...):正式使用进阶特性/Scale 版时配置授权密钥;仅用 Community 版可注释掉该行;- 启动后
time.sleep(20):等待 Kafka 完全就绪再pw.run(),避免容器编排时序导致连接失败。
验证输出:读回结果并落地为 CSV
统一后的日志已经出现在 Kafka 的 unified_timestamps topic 中,可以使用任意 Kafka 客户端读取。更省事的做法是继续用 Pathway 检查一切是否正常:在 pathway-src/ 下创建 read-results.py 读取该 topic 并写成 CSV:
table = pw.io.kafka.read(
rdkafka_settings,
topic=topic_name,
schema=InputStreamSchema,
format="json",
autocommit_duration_ms=100,
)
pw.io.csv.write(table, "./results.csv")
pw.run()
仓库中的 read-results.py 与上述逻辑一致,并额外用 str(uuid4()) 生成不重复的消费组 ID,避免与 ETL 主进程的位点互相干扰;读取端的 InputStreamSchema 改为:
class InputStreamSchema(pw.Schema):
timestamp: float
message: str
运行方式可参考 README.md:make connect(或 docker compose exec -it pathway bash)进入 pathway 容器后执行 python read-results.py,Pathway 会持续读取并追加写入 results.csv。脚本产出的结果形如:
timestamp,message,time,diff
1707217879632.242,"11",1707217879944,1
1707217876629.236,"8",1707217879944,1
1707217872469.24,"4",1707217879944,1
1707217868355.006,"0",1707217879944,1
1707217870466.797,"2",1707217879944,1
1707217873626.241,"5",1707217879944,1
1707217869465.5308,"1",1707217879944,1
1707217871468.065,"3",1707217879944,1
1707217874627.24,"6",1707217879944,1
1707217877630.239,"9",1707217879944,1
1707217875628.488,"7",1707217879944,1
1707217878631.242,"10",1707217879944,1
1707217880633.24,"12",1707217880644,1
1707217881634.5,"13",1707217881644,1
1707217882635.752,"14",1707217882644,1
可以看到:date/message 两列已被加工成 timestamp 毫秒时间戳列 + message 编号列,纽约与巴黎两个时区的日志时间都被统一成了同一时间轴上的绝对时刻——时区转换的目标达成。
完整方案:一文件全量代码
下面给出完整的 etl.py,将抽取、转换、装载串成一条管道(以仓库当前版本为准,去除了授权密钥行后的核心逻辑):
import time
import pathway as pw
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
timestamps_timezone_1 = pw.io.kafka.read(
rdkafka_settings,
topic="timezone1",
format="json",
schema=InputStreamSchema,
autocommit_duration_ms=100,
)
timestamps_timezone_2 = pw.io.kafka.read(
rdkafka_settings,
topic="timezone2",
format="json",
schema=InputStreamSchema,
autocommit_duration_ms=100,
)
def convert_to_timestamp(table):
table = table.select(
date=pw.this.date.dt.strptime(fmt=str_repr, contains_timezone=True),
message=pw.this.message,
)
table_timestamp = table.select(
timestamp=pw.this.date.dt.timestamp(unit="ms"),
message=pw.this.message,
)
return table_timestamp
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)
pw.io.kafka.write(
timestamps_unified, rdkafka_settings, topic_name="unified_timestamps", format="json"
)
# 等 Kafka 就绪后启动计算
time.sleep(20)
pw.run()
延伸阅读与可继续深入的方向
至此,你已经能用 Pathway 完成一个标准的 Kafka ETL 管道。实际业务里你的环境可能与演示略有差异,需要不同的算子组合。可以继续从以下方向深入,素材都沉淀在当前仓库中:
- 连接器体系:除 Kafka 外,Pathway 还提供大量输入/输出连接器(CSV、文件系统、各类数据库、消息队列等)用于抽取与装载,输入侧可在
python/pathway/io/下浏览实现,例如 Kafka 相关代码集中在 python/pathway/io/kafka/init.py; - 表算子与数据转换:本文用到的
select、concat_reindex只是基础表操作的一部分,更多拼接、过滤、分组、排序等算子可在 python/pathway/internals/table.py 中按需查阅; - 时间态(temporal)计算:若后续需要在统一时间轴上做跨流对齐(例如把告警日志与交易流水按时间关联),可以研究 ASOF join、interval join 等时间对齐算子;
- 工程复现:整个可运行的演示工程在 examples/projects/kafka-ETL,结合 docker-compose.yml 与 Makefile 即可一键起停、查看日志并核对输出。
一句话总结方法论:用 Schema 声明结构、用 pw.io.kafka.read 抽取、用 dt 命名空间把带时区字符串解析为绝对时间戳、用 concat_reindex 合并、用 pw.io.kafka.write 装载、最后 pw.run() 让数据真正开始流动。
atomcodeClaude Code 的开源替代方案。连接任意大模型,编辑代码,运行命令,自动验证 — 全自动执行。用 Rust 构建,极致性能。 | An open-source alternative to Claude Code. Connect any LLM, edit code, run commands, and verify changes — autonomously. Built in Rust for speed. Get StartedRust0626
Hy4-previewHy4 preview 是由腾讯混元团队研发的新一代混合专家(MoE)旗舰模型。模型总参数量 770B,每个 token 激活 49B,主干共包含78层,第一层采用标准 FFN,其余 77 层均为 MoE 结构,每层包含 256 个路由专家与 1 个共享专家,每个 token 激活 top-8 路由专家及共享专家。主干之外原生内置 1 层 MTP(总参数量 10B,激活 0.7B)以支持投机解码。Python00
GLM-5.3GLM-5.3 与 GLM-5.2 使用相同的基座模型——所有提升均来自后训练。与 GLM-5.2 相比,它在复杂编程和长程任务上的表现显著提升。Jinja00
GLM-5.3-FlashGLM-5.3-Flash (320B-A18B),是GLM-5系列的首个原生多模态模型。320B总参数,能力超过GLM-5.2Jinja00
Spark-X2.5-4BSpark-X2.5-4B 旨在让强大的 AI 更实用、更高效、更易获得。在广泛日常任务中表现强劲,涵盖对话、写作、翻译、推理、编码、工具调用以及智能体工作流,并在同等规模的开源模型中取得领先成绩。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00
Spark-X2.5-1.7BSpark-X2.5-1.7B 旨在让强大的 AI 更加实用、高效且易于获取。这些模型在广泛的日常任务中表现出色,涵盖对话、写作、翻译、推理、编程、工具调用和智能体工作流,并在同等规模的开源模型中取得领先结果。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00