首页
/ Pathway 实现 Kafka ETL:用 Python 统一跨时区事件流并写回 Kafka

Pathway 实现 Kafka ETL:用 Python 统一跨时区事件流并写回 Kafka

2026-09-07 11:41:54作者:钟日瑜

本文以 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:9092KAFKA_AUTO_CREATE_TOPICS: true(允许 topic 自动创建,示例因此无需手工建 topic)等环境变量,并把 9092 映射到宿主机;
  • pathwaystream-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 builddocker compose up -d 拉起全部容器;
  • make stopdocker 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.dumpsutf-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.serversgroup.id(见 python/pathway/io/kafka/init.py),并支持 start_from_timestamp_msparallel_readerswith_metadatajson_field_paths 等进阶参数,用于指定消费起点、并行读、读取消息元数据或按 JSON Pointer 重映射字段。

由于 Kafka 消息中 date 是形如 2024-02-05 10:01:52.884548 -0500 的字符串,这里 format="json" + schema 会直接把 datemessage 两个 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.serverstopic_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.mdmake 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
  • 表算子与数据转换:本文用到的 selectconcat_reindex 只是基础表操作的一部分,更多拼接、过滤、分组、排序等算子可在 python/pathway/internals/table.py 中按需查阅;
  • 时间态(temporal)计算:若后续需要在统一时间轴上做跨流对齐(例如把告警日志与交易流水按时间关联),可以研究 ASOF join、interval join 等时间对齐算子;
  • 工程复现:整个可运行的演示工程在 examples/projects/kafka-ETL,结合 docker-compose.ymlMakefile 即可一键起停、查看日志并核对输出。

一句话总结方法论:用 Schema 声明结构、用 pw.io.kafka.read 抽取、用 dt 命名空间把带时区字符串解析为绝对时间戳、用 concat_reindex 合并、用 pw.io.kafka.write 装载、最后 pw.run() 让数据真正开始流动。

登录后查看全文
热门项目推荐
相关项目推荐