首页
/ Pathway 从 Kafka 迁移到 Redpanda:零代码改动的流处理实战指南

Pathway 从 Kafka 迁移到 Redpanda:零代码改动的流处理实战指南

2026-09-06 21:05:04作者:范靓好Udolf

本篇基于 Pathway(Python ETL / 实时流处理框架)官方教程《Switching from Kafka to Redpanda》,带你以“最佳评分电影 Top-K”为例,完整走通 Kafka 版项目部署、切换到 Redpanda 的全部改动点,以及一个必须注意的 topic 自动创建丢消息陷阱;读完后你将掌握:如何在 Pathway 中用 pw.io.kafka / pw.io.redpanda 连接器接入两类消息流平台、如何用 Docker Compose 一键编排 Redpanda,以及从源码层面理解“为什么切换几乎不需要改代码”。

背景:为什么有人想从 Kafka 换到 Redpanda

Kafka 在 2011 年发布时是实时数据处理领域的变革者:它提供的分布式流式平台让组织可以构建可扩展、容错的数据管道,实时处理数据流并接入各类数据源与数据汇。

但随着生态演进,Kafka 的一些当年领先特性逐渐被部分开发者视为负担:部署与管理相对复杂、依赖 ZooKeeper、运行在 JVM 上。Redpanda 正是针对这些痛点出现的方案:它是一个用 C++ 重写、宣称 Kafka 完全兼容的高性能流平台,定位为更低的延迟、更高效的资源利用率和更简单的部署运维。

那么“Kafka 兼容”到底有多好迁移?本文的核心问题就是:把基于 Pathway Live Data Framework 的 Kafka 集成项目改到 Redpanda 上,到底要付出多少代价?

先给出结论(教程原文的“剧透”):Pathway Live Data Framework 对 Kafka 和 Redpanda 完全透明,你的业务代码一行都不用改。

这个结论不是空话,仓库源码可以直接印证。Redpanda 连接器实现 中,read() 函数在参数校验后直接委托给 Kafka 连接器:

# python/pathway/io/redpanda/__init__.py(节选)
return kafka.read(
    rdkafka_settings=rdkafka_settings,
    topic=topic,
    schema=schema,
    mode=mode,
    format=format,
    schema_registry_settings=schema_registry_settings,
    debug_data=debug_data,
    autocommit_duration_ms=autocommit_duration_ms,
    json_field_paths=json_field_paths,
    parallel_readers=parallel_readers,
    name=name,
    max_backlog_size=max_backlog_size,
    _stacklevel=5,
)

write() 同样是直接转发到 kafka.write(...)(见 python/pathway/io/redpanda/init.py)。也就是说,从源码结构看,pw.io.redpandapw.io.kafka 的一层“品牌别名”,底层共用同一套 librdkafka 连接逻辑——而 Redpanda 本身兼容 Kafka API,所以“换平台不改代码”成立。

示例场景:找出评分最高的 K 部电影

教程用一个具体的业务问题贯穿全文:你刚入职一个流行的 VOD 流媒体平台,第一个任务是找出目录中评分最高的 K 部电影,以及每部电影收到了多少评分。例如 K=3 时期望的输出表:

MovieID Average RatingNumber
1 218 4.9 7510
2 45 4.8 9123
3 7456 4.8 1240

评分数据以流的形式经由 Kafka 实例到达,你需要把实时更新的排行榜输出到 CSV 文件。Pathway 侧负责“监听 topic → 流式计算 → 持续写出”,整个流程天然契合 Pathway 的流式计算模型:每当新评分到达,Top-K 结果自动更新,而不是全量重算。

仓库中提供了完整可运行的示例工程,Kafka 版与 Redpanda 版分别位于 examples/projects/best-movies-example/kafka-versionexamples/projects/best-movies-example/redpanda-version,两版目录结构一致:

.
├── pathway-src/
│   ├── Dockerfile
│   └── process-stream.py
├── producer-src/
│   ├── create-stream.py
│   ├── dataset.csv
│   └── Dockerfile
├── docker-compose.yml
└── Makefile

Kafka 版方案:四个组件搭起来

组件与职责

Kafka 方案需要四个组件,各自独立 Docker 容器:

  • ZooKeeper
  • Kafka
  • Pathway Live Data Framework
  • 流生产者(stream producer)

流生产者把评分发送到 Kafka 的 ratings topic;Pathway 监听该 topic,处理流并输出排行榜到 best_rating.csv

配置 ZooKeeper 与 Kafka

两者都在 docker-compose.yml 中配置。教程为保持简洁未启用任何安全机制,完整配置如下(对应仓库文件 kafka-version/docker-compose.yml):

version: "3.7"
name: tuto-switch-to-redpanda
networks:
  tutorial_network:
    driver: bridge
services:
  zookeeper:
    image: confluentinc/cp-zookeeper:5.5.3
    environment:
      ZOOKEEPER_CLIENT_PORT: 2181
    networks:
      - tutorial_network
  kafka:
    image: confluentinc/cp-enterprise-kafka:5.5.3
    depends_on: [zookeeper]
    environment:
      KAFKA_AUTO_CREATE_TOPICS: true
      KAFKA_ZOOKEEPER_CONNECT: "zookeeper:2181"
      KAFKA_ADVERTISED_HOST_NAME: kafka
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9092
      KAFKA_BROKER_ID: 1
      KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
      KAFKA_JMX_PORT: 9991
      KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT
      KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT
      CONFLUENT_SUPPORT_METRICS_ENABLE: false
    ports:
    - 9092:9092
    command: sh -c "((sleep 15 && kafka-topics --create --zookeeper zookeeper:2181 --replication-factor 1 --partitions 1 --topic ratings)&) && /etc/confluent/docker/run "
    networks:
      - tutorial_network

几个关键点:

  • 消息发往名为 ratings 的 topic,它在 command 中延迟 15 秒后用 kafka-topics --create 显式创建(单分区、副本因子 1);
  • KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9092 决定了容器内其他服务必须用 kafka:9092(而不是 localhost)来连接;
  • 未启用 SASL/SSL,属于纯教程环境,生产环境需另行配置安全机制。

生成数据流

流来自一个 CSV 数据集,列定义为:userId (int)、movieId (int)、rating (float)、timestamp (int)——这正是 GroupLens 的 MovieLens25M 数据集所用的 schema。教程提供了一个玩具数据集作演示,但同一套项目可以直接跑完整 MovieLens25M:

userId,movieId,rating,timestamp
1,296,5.0,1147880044
1,306,3.5,1147868817
1,307,5.0,1147868828
1,665,5.0,1147878820
1,899,3.5,1147868510
1,1088,4.0,1147868495
2,296,4.0,1147880044
2,306,2.5,1147868817
2,307,3.0,1147868828
2,665,2.0,1147878820
2,899,4.5,1147868510
2,1088,2.0,1147868495
3,296,1.0,1147880044
3,306,2.5,1147868817
3,307,4.0,1147868828
3,665,2.0,1147878820
3,899,1.5,1147868510
3,1088,5.0,1147868495

生成流的脚本逐行读取 CSV,把每条评分序列化为 JSON 消息后发送给 Kafka(完整脚本见 producer-src/create-stream.py):

import csv
import json
import time

from kafka import KafkaProducer

topic = "ratings"

# 等待 Kafka 和 ZooKeeper 就绪
time.sleep(30)

producer = KafkaProducer(
    bootstrap_servers=["kafka:9092"],
    security_protocol="PLAINTEXT",
    api_version=(0, 10, 2),
)

with open("./dataset.csv", newline="") as csvfile:
    dataset_reader = csv.reader(csvfile, delimiter=",")
    first_line = True
    for row in dataset_reader:
        # 跳过表头
        if first_line:
            first_line = False
            continue
        message_json = {
            "userId": int(row[0]),
            "movieId": int(row[1]),
            "rating": float(row[2]),
            "timestamp": int(row[3]),
        }
        producer.send(topic, (json.dumps(message_json)).encode("utf-8"))
        time.sleep(0.1)

producer.send(topic, "*COMMIT*".encode("utf-8"))
time.sleep(2)
producer.close()

注意这里连接的是 kafka:9092 而非 localhost,因为生产者运行在同一个 Docker 网络内。该脚本运行在自己的容器里:

  stream-producer:
    build:
      context: .
      dockerfile: ./producer-src/Dockerfile
    depends_on: [kafka]
    networks:
      - tutorial_network

Dockerfile 也只需 Python 镜像加 kafka-python 包:

FROM python:3.10

RUN pip install kafka-python
COPY ./producer-src/create-stream.py create-stream.py
COPY ./producer-src/dataset.csv dataset.csv

CMD ["python", "-u", "create-stream.py"]

Pathway 侧:连接 Kafka 并计算 Top-K

Pathway 连接 Kafka 的核心是 rdkafka_settings 字典(即 librdkafka 配置格式)。明文连接:

rdkafka_settings = {
    "bootstrap.servers": "kafka:9092",
    "security.protocol": "plaintext",
    "group.id": "0",
    "session.timeout.ms": "6000",
}

如果需要更安全的连接,可以使用 SASL-SSL + SCRAM-SHA-256 机制:

rdkafka_settings = {
    "bootstrap.servers": "server:9092",
    "security.protocol": "sasl_ssl",
    "sasl.mechanism": "SCRAM-SHA-256",
    "group.id": "$GROUP_NAME",
    "session.timeout.ms": "6000",
    "sasl.username": "username",
    "sasl.password": "********"
}

然后用 pw.io.kafka.read 连接器读取 ratings topic:

class InputSchema(pw.Schema):
    movieId: int
    rating: float

t_ratings = pw.io.kafka.read(
    rdkafka_settings,
    topic="ratings",
    format="json",
    schema=InputSchema,
    autocommit_duration_ms=100,
)

两点说明(结合 连接器源码 的参数签名,Kafka 连接器同构):

  • schema 只声明需要的列。消息里虽然带 userIdtimestamp,但我们只关心 movieIdrating,schema 里不需要包含其他字段;
  • autocommit_duration_ms 控制提交节奏。它是两次提交之间的最大间隔,连接器会把这段时间内收到的更新打包提交进 Pathway 计算图。连接器默认值为 1500 毫秒,教程显式设为 100 毫秒以获得更灵敏的批处理节奏。

Top-K 计算函数 compute_best 完整实现如下:先按 movieId 分组聚合求和与计数,再算平均分,然后把“平均分、评分数、movieId”打成一个 tuple,用 sorted_tuple 维护全局有序集合,取末尾 K 个并展平回列:

def compute_best(t_ratings, K):
    t_best_ratings = t_ratings.groupby(pw.this.movieId).reduce(
        pw.this.movieId,
        sum_ratings=pw.reducers.sum(pw.this.rating),
        number_ratings=pw.reducers.count(pw.this.rating),
    )
    t_best_ratings = t_best_ratings.select(
        pw.this.movieId,
        pw.this.number_ratings,
        average_rating=pw.apply(
            lambda x, y: (x / y) if y != 0 else 0,
            pw.this.sum_ratings,
            pw.this.number_ratings,
        ),
    )
    t_best_ratings = t_best_ratings.select(
        movie_tuple=pw.apply(
            lambda x, y, z: (x, y, z),
            pw.this.average_rating,
            pw.this.number_ratings,
            pw.this.movieId,
        )
    )
    t_best_ratings = t_best_ratings.reduce(
        total_tuple=pw.reducers.sorted_tuple(pw.this.movie_tuple)
    )
    t_best_ratings = t_best_ratings.select(
        K_best=pw.apply(lambda my_tuple: (list(my_tuple))[-K:], pw.this.total_tuple)
    )
    t_best_ratings = t_best_ratings.flatten(pw.this.K_best).select(
        pw.this.K_best
    )
    t_best_ratings = t_best_ratings.select(
        movieId=pw.apply(lambda rating_tuple: rating_tuple[2], pw.this.K_best),
        average_rating=pw.apply(lambda rating_tuple: rating_tuple[0], pw.this.K_best),
        views=pw.apply(lambda rating_tuple: rating_tuple[1], pw.this.K_best),
    )
    return t_best_ratings

最终的完整处理文件(Kafka 版):

import pathway as pw
import time

rdkafka_settings = {
    "bootstrap.servers": "kafka:9092",
    "security.protocol": "plaintext",
    "group.id": "0",
    "session.timeout.ms": "6000",
}

class InputSchema(pw.Schema):
    movieId: int
    rating: float

t_ratings = pw.io.kafka.read(
    rdkafka_settings,
    topic="ratings",
    format="json",
    schema=InputSchema,
    autocommit_duration_ms=100,
)

t_best_ratings = compute_best(t_ratings, 3)

# 结果输出到专门的 CSV 文件
pw.io.csv.write(t_best_ratings, "./best_ratings.csv")

# 等待 Kafka 和 ZooKeeper 就绪
time.sleep(20)
# 启动计算
pw.run()

Pathway 也有自己的容器:

  pathway:
    build:
      context: .
      dockerfile: ./pathway-src/Dockerfile
    depends_on: [kafka]
    networks:
      - tutorial_network
FROM python:3.10

RUN pip install -U pathway
COPY ./pathway-src/process-stream.py process-stream.py

CMD ["python", "-u", "process-stream.py"]

运行结果

用上面的玩具数据集运行,输出 CSV(带 time/diff 等流式输出元数据列)形如:

movieId,average_rating,views,time,diff
296,5,1,1680008702067,1
306,3.5,1,1680008702167,1
[...]
296,3.3333333333333335,3,1680008703767,-1
1088,3.6666666666666665,3,1680008703767,1
296,3.3333333333333335,3,1680008703767,1
899,3.1666666666666665,3,1680008703767,-1

可以看到,每当新评分到达导致排名变化时,Top-3 会被持续更新——这正是流式 Top-K 而非一次性批处理的表现。

切换到 Redpanda:改动点逐项过一遍

现在团队引入了 Redpanda:它宣称完全兼容 Kafka API、不依赖 ZooKeeper、更易管理,并且由于没有 page cache 数据落盘更持久,其公开基准测试的吞吐也优于 Kafka。任务是:把项目从 Kafka 切到 Redpanda。

切换到 Redpanda 后,项目组件从 4 个变 3 个:Redpanda、Pathway、流生产者(ZooKeeper 直接消失)。

第一步:Docker Compose 里用 redpanda 服务替换 kafka + zookeeper

若已有现成 Redpanda 实例可跳过此节。删掉 kafkazookeeper 两个服务,换成一个 redpanda 服务(对应 redpanda-version/docker-compose.yml):

services:
  redpanda:
    command:
      - redpanda
      - start
      - --kafka-addr internal://0.0.0.0:9092,external://0.0.0.0:19092
      - --advertise-kafka-addr internal://redpanda:9092,external://localhost:19092
      - --pandaproxy-addr internal://0.0.0.0:8082,external://0.0.0.0:18082
      - --advertise-pandaproxy-addr internal://redpanda:8082,external://localhost:18082
      - --schema-registry-addr internal://0.0.0.0:8081,external://0.0.0.0:18081
      - --rpc-addr redpanda:33145
      - --advertise-rpc-addr redpanda:33145
      - --smp 1
      - --memory 1G
      - --mode dev-container
      - --default-log-level=debug
      - --set redpanda.enable_transactions=true
      - --set redpanda.enable_idempotence=true
      - --set redpanda.auto_create_topics_enabled=true
    image: docker.redpanda.com/redpandadata/redpanda:v23.1.2
    container_name: redpanda
    volumes:
      - redpanda:/var/lib/redpanda/data
    networks:
      - tutorial_network

参数逐项解释:

参数 作用
--kafka-addr / --advertise-kafka-addr 定义 Kafka 协议监听与对外宣告地址。internal://redpanda:9092 供同网络容器连接,external://localhost:19092 供宿主机访问
--pandaproxy-addr / --schema-registry-addr HTTP 代理(8082)与 Schema Registry(8081)的地址,本教程用不到但示例中一并开启
--rpc-addr / --advertise-rpc-addr 集群内部节点间通信端口(33145)
--smp 1--memory 1G 单核、1GB 内存,典型的开发容器资源限制
--mode dev-container 开发容器模式,放宽单节点部署校验
--set redpanda.enable_transactions=true / enable_idempotence=true 开启事务与幂等生产,对齐 Kafka 语义
--set redpanda.auto_create_topics_enabled=true 允许首条消息触发自建 topic——这一项直接关联后面要讲的丢消息陷阱

注意容器名变化:kafka 改成了 redpanda,所以客户端的 bootstrap 地址需要从 kafka:9092 改为 redpanda:9092。教程也给了一个偷懒技巧:把 Redpanda 容器命名为 kafka(或反过来改 Kafka 容器名),连地址都不用动。

第二步:Pathway 侧——只改服务名

process-stream.py 里唯一的实质改动是 bootstrap.servers

rdkafka_settings = {
    "bootstrap.servers": "redpanda:9092",
    "security.protocol": "plaintext",
    "group.id": "0",
    "session.timeout.ms": "6000",
}

配置文件其余内容与 Kafka 时完全一致。为了“名实相符”,Pathway 也提供了专门的 Redpanda 连接器 pw.io.redpanda.read,其签名与 pw.io.kafka.read 完全同构,完整文件如下(与仓库 redpanda-version/pathway-src/process-stream.py 基本一致,仓库文件额外包含一行 pw.set_license_key,聚合处写作 pw.reducers.count()):

import pathway as pw
import time

rdkafka_settings = {
    "bootstrap.servers": "redpanda:9092",
    "security.protocol": "plaintext",
    "group.id": "0",
    "session.timeout.ms": "6000",
}

class InputSchema(pw.Schema):
    movieId: int
    rating: float

t_ratings = pw.io.redpanda.read(
    rdkafka_settings,
    topic="ratings",
    format="json",
    schema=InputSchema,
    autocommit_duration_ms=100,
)

t_best_ratings = compute_best(t_ratings, 3)

# 结果输出到专门的 CSV 文件
pw.io.csv.write(t_best_ratings, "./best_ratings.csv")

# 等待 Redpanda 就绪
time.sleep(20)

# 启动计算
pw.run()

如果你不在意连接器和服务的名字,这个文件可以完全不动——pw.io.kafka.read 指向 Redpanda 一样工作。

第三步:流生产者侧 + 一个 Redpanda 特有的坑

生产者脚本同样只需改服务名:

producer = KafkaProducer(
    bootstrap_servers=["redpanda:9092"],
    security_protocol="PLAINTEXT",
    api_version=(0, 10, 2),
)

但此时运行会发现 best_ratings.csv空的。原因是两边 topic 创建时机不同:

  • Kafka 版中,topic 已在 command 里提前创建,消费者一开始就能就位;
  • Redpanda 中 topic 由第一条消息触发自动创建,而创建期间 Redpanda 会丢弃期间到达的消息。

解决办法:在流式发送数据之前,先发送一条消息“预热”,让 topic 提前进入就绪状态:

producer.send(topic, "*COMMIT*".encode("utf-8"))
time.sleep(2)

教程特别强调:这个差异来自 Redpanda 本身的消息处理策略,与 Pathway 无关——Pathway 对两者完全透明,消息的丢弃发生在 broker 侧。修复后的生产者完整脚本(与仓库 redpanda-version/producer-src/create-stream.py 一致):

from kafka import KafkaProducer
import csv
import time
import json

topic = "ratings"

# 等待 Redpanda 就绪
time.sleep(30)

producer = KafkaProducer(
    bootstrap_servers=["redpanda:9092"],
    security_protocol="PLAINTEXT",
    api_version=(0, 10, 2),
)
producer.send(topic, "*COMMIT*".encode("utf-8"))
time.sleep(2)

with open("./dataset.csv", newline="") as csvfile:
    dataset_reader = csv.reader(csvfile, delimiter=",")
    first_line = True
    for row in dataset_reader:
        # 跳过表头
        if first_line:
            first_line = False
            continue
        message_json = {
            "userId": int(row[0]),
            "movieId": int(row[1]),
            "rating": float(row[2]),
            "timestamp": int(row[3]),
        }
        producer.send(topic, (json.dumps(message_json)).encode("utf-8"))
        time.sleep(0.1)

producer.send(topic, "*COMMIT*".encode("utf-8"))
time.sleep(2)
producer.close()

用 Makefile 运行与验证结果

Redpanda 版工程附带 Makefile,常用目标:

build:
	docker compose up -d
stop:
	docker compose down -v
	docker rmi best-movies-redpanda-pathway:latest
	docker rmi best-movies-redpanda-stream-producer:latest
logs:
	docker compose logs pathway
logs-prod:
	docker compose logs stream-producer

make build 起服务、make stop 清理容器与镜像、make logs 查看 Pathway 容器日志;另有 rpk cluster inforpk topic create/produce 等目标用于直接操作 Redpanda 集群。

修复预热消息后,最终输出与 Kafka 版完全相同:

movieId,average_rating,views,time,diff
296,5,1,1680008702067,1
306,3.5,1,1680008702167,1
[...]
296,3.3333333333333335,3,1680008703767,-1
1088,3.6666666666666665,3,1680008703767,1
296,3.3333333333333335,3,1680008703767,1
899,3.1666666666666665,3,1680008703767,-1

进阶:把结果写回 Redpanda

写出到本地 CSV 有两个局限:CSV 不是流式数据的理想载体,且文件留在本地,组织内其他人无法实时消费。更好的做法是把排行榜实时写回 Redpanda 的 best_ratings topic——Redpanda 天然为流式消费优化。

pw.io.redpanda.write 连接器即可(与 Kafka 写连接器完全等价):

rdkafka_settings = {
    "bootstrap.servers": "redpanda:9092",
    "security.protocol": "plaintext",
    "group.id": "$GROUP_NAME",
    "session.timeout.ms": "6000",
}
pw.io.redpanda.write(
    t_best_ratings,
    rdkafka_settings,
    topic_name="best_ratings",
    format="json"
)

连接 Redpanda Cloud 等托管实例时,同样支持 SASL-SSL + SCRAM-SHA-256:

rdkafka_settings = {
    "bootstrap.servers": "redpanda:9092",
    "security.protocol": "sasl_ssl",
    "sasl.mechanism": "SCRAM-SHA-256",
    "group.id": "$GROUP_NAME",
    "session.timeout.ms": "6000",
    "sasl.username": "username",
    "sasl.password": "********",
}

结合 write 连接器源码 还可以补充几点实操信息:

  • format 目前仅支持 "json"(源码注释说明这是为将来扩展预留的可配置项);
  • topic_name 外还支持 subject(配合 Confluent Schema Registry)、name(用于日志、监控面板及持久化快照命名)、sort_by(对每个 minibatch 内的输出按指定列升序排序)等参数;
  • 表的所有更新都会被发送到该 topic,下游任何消费者都能实时获取最新排行榜。

总结

  • 迁移代价趋近于零:Redpanda 完全兼容 Kafka API,而 Pathway 的 Redpanda 连接器在源码层面就是 Kafka 连接器的直接转发(见 python/pathway/io/redpanda/init.py),因此若容器名保持不变,业务代码可以完全不动;在意命名一致性时,把 pw.io.kafka 换成 pw.io.redpanda 即可。
  • 架构简化:删掉 ZooKeeper,四组件变三组件,Redpanda 单容器即可承载。
  • 唯一需要警惕的坑:Redpanda 的 topic 由首条消息自动创建,创建期间到达的消息会被 broker 丢弃,需要一条“预热”消息(如 *COMMIT*)先行触发 topic 创建;这是 Redpanda 侧行为,与 Pathway 无关。
  • 结果闭环:算好的 Top-K 排行榜既可以用 pw.io.csv.write 落盘,也可以用 pw.io.redpanda.write 实时写回 best_ratings topic 供全组织消费。

完整可运行工程可直接参考仓库中的 Kafka 版示例Redpanda 版示例

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