PySpark Kafka 集成测试指南:用 KafkaUtils 与 Docker 容器在本地测试 Kafka 流处理
2026/9/20 11:04:22 网站建设 项目流程
  • 大数据
  • 数据分析
  • 批处理
  • 流处理
  • 机器学习
  • 图计算

【免费下载链接】spark

Apache Spark - A unified analytics engine for large-scale data processing

项目地址:https://gitcode.com/gh_mirrors/sp/spark
点击查看免费下载

导读

本指南面向需要在本地对 PySpark Kafka 流处理应用做端到端集成测试的开发者。文档以仓库中 python/pyspark/sql/tests/streaming/KAFKA_TESTING.md 为骨架,结合 kafka_utils.py 与 test_streaming_kafka_rtm.py 的源码实现,完整讲解如何通过 testcontainers-python 在 Docker 中拉起单 broker Kafka 集群,并用KafkaUtils完成建 Topic、发消息、Spark 读写、流式查询断言等全流程测试。读完后你将能独立编写可复制、可运行、可反复执行的 Kafka 集成测试用例,并理解其底层原理。

为什么需要 Docker 化的 Kafka 测试

PySpark 的 Kafka 集成测试属于典型的端到端测试:既要有真实可用的 Kafka broker,又要验证 Spark 的 Kafka 数据源(format("kafka"))能否正确读写。手工搭建 Kafka 集群繁琐且难以清理,而 mock 又无法覆盖真实的序列化、分区、偏移量提交等行为。

仓库的方案是:在测试代码中直接通过 Docker 启动一个单 broker 的 Confluent Kafka 容器,全部生命周期由KafkaUtils封装管理。这样:

  • 零手工配置:不需要本地安装 Kafka、不需要手写 server.properties、不需要管理 Zookeeper/KRaft 元数据;
  • 可重复、可清理setup()启动容器,teardown()关闭容器并清理客户端资源,测试结束不留残留进程;
  • 贴近生产:测试的是真实 broker,Spark 读到的 key/value 是真实的二进制数据,行为与生产环境一致。

从源码结构看,这一工具的设计意图是与 Python unittest 体系(尤其是ReusedSQLTestCase)深度配合,作为整个测试类的类级夹具使用。

环境准备

1. Docker

KafkaUtils依赖 testcontainers 调用本机 Docker daemon 启动容器,因此 Docker 必须已安装并处于运行状态。验证方式:

docker ps

仓库中的测试同样在类装饰器层面做了一层 Docker 可用性探测(见 test_streaming_kafka_rtm.py):通过docker info命令探测 daemon,不可用时整类测试被跳过,避免 CI 环境无 Docker 时误报失败。

2. Python 依赖

安装运行测试所需的两个核心包:

pip install testcontainers[kafka] kafka-python-ng

也可以按项目方式安装全部开发依赖:

cd $SPARK_HOME pip install --group dev

setup()源码中(kafka_utils.py)对依赖做了显式校验:缺少testcontainers.kafkakafka(KafkaProducer/KafkaAdminClient)都会抛出带安装提示的ImportError。测试端对应的跳检逻辑在 python/pyspark/testing/utils.py,通过have_package探测kafkatestcontainers,缺失时由@unittest.skipIf跳过测试类。

3. Spark 构建产物

运行 Kafka 测试还需要 Spark 的 Kafka SQL 连接器 JAR。测试基类StreamingKafkaTestsMixin在创建 SparkSession 之前会通过search_jarconnector/kafka-0-10-sql目录中定位spark-sql-kafka-0-10_的 JAR,并把全部依赖写入PYSPARK_SUBMIT_ARGS --jars(见 test_streaming_kafka_rtm.py)。若 JAR 不存在,会提示先构建 Spark:

build/mvn package # 或 build/sbt Test/package

快速开始:第一个 Kafka 测试

最小可用用例结构如下:

import unittest from pyspark.sql.tests.streaming.kafka_utils import KafkaUtils from pyspark.testing.sqlutils import ReusedSQLTestCase class MyKafkaTest(ReusedSQLTestCase): @classmethod def setUpClass(cls): super().setUpClass() cls.kafka_utils = KafkaUtils() cls.kafka_utils.setup() @classmethod def tearDownClass(cls): cls.kafka_utils.teardown() super().tearDownClass() def test_kafka_read_write(self): # Create a topic topic = "test-topic" self.kafka_utils.create_topics([topic]) # Send test data messages = [("key1", "value1"), ("key2", "value2")] self.kafka_utils.send_messages(topic, messages) # Read with Spark df = ( self.spark.read .format("kafka") .option("kafka.bootstrap.servers", self.kafka_utils.broker) .option("subscribe", topic) .option("startingOffsets", "earliest") .load() ) # Verify data results = df.selectExpr( "CAST(key AS STRING) as key", "CAST(value AS STRING) as value" ).collect() self.assertEqual(len(results), 2)

几个关键点:

  • 测试类继承ReusedSQLTestCase(定义于 python/pyspark/testing/sqlutils.py),提供self.spark(复用 JVM 与 SparkSession)及 SQL 断言工具;
  • Kafka 容器在setUpClass类级启动、类级关闭,一个测试类只启停一次,避免每个用例都付出 10~30 秒的容器拉起成本;
  • broker 地址通过self.kafka_utils.broker动态获取,每次运行时端口随机分配,不写死端口;
  • Kafka 数据源的 key/value 是二进制列,读取时用CAST(... AS STRING)反序列化。

KafkaUtils API 参考(基于源码逐项解析)

以下 API 说明与 kafka_utils.py 源码一一对应。

初始化与生命周期

__init__(kafka_version="7.4.0")

创建一个KafkaUtils实例。

  • kafka_version:使用的 Confluent Kafka 镜像版本,默认"7.4.0"。源码在setup()中将其拼为confluentinc/cp-kafka:{kafka_version}传给KafkaContainer(kafka_utils.py)。默认选 7.4.0 是出于稳定性考量。
setup()

启动 Kafka 容器并初始化 admin client 与 producer。必须在调用任何其他方法之前调用。

内部流程(kafka_utils.py):

  1. 幂等检查:若initialized为 True 直接返回;
  2. 导入testcontainers.kafka.KafkaContainer,缺失则抛ImportError
  3. 导入kafka.KafkaProducerkafka.admin.KafkaAdminClient,缺失则抛ImportError
  4. 启动容器,通过get_bootstrap_server()获得broker地址(形如localhost:9093,端口由 testcontainers 随机映射);
  5. 用 10 秒超时初始化KafkaAdminClientKafkaProducer(producer 的 key/value 序列化器会把非 None 值str()后 UTF-8 编码);
  6. 任意一步异常都会先调用teardown()清理再重新抛出RuntimeError

可能抛出的异常

  • ImportError:依赖未安装;
  • RuntimeError:容器启动失败。

注意:首次运行时 Docker 需要拉取镜像,容器启动可能耗时 10~30 秒。

teardown()

关闭 admin client、关闭 producer(5 秒超时)、停止容器并复位所有状态字段。实现上每一步都做了异常兜底与字段置空(kafka_utils.py),因此可安全多次调用,适合放在finallytearDownClass中。

内部防护:_assert_initialized()

所有业务方法(建 Topic、删 Topic、发消息、读记录、访问属性)执行前都会调用_assert_initialized();未调用setup()时抛出RuntimeError("KafkaUtils has not been initialized. Call setup() first.")(kafka_utils.py)。

Topic 管理

create_topics(topic_names, num_partitions=1, replication_factor=1)

批量创建 Topic。

  • topic_namesList[str]):要创建的 Topic 名列表;
  • num_partitions(int):每个 Topic 的分区数,默认 1;
  • replication_factor(int):副本因子,默认 1;单 broker 环境下最大只能是 1。

实现上通过KafkaAdminClient.create_topics提交NewTopic,若TopicAlreadyExistsError则静默跳过(kafka_utils.py)。

# Create single partition topics kafka_utils.create_topics(["topic1", "topic2"]) # Create multi-partition topic kafka_utils.create_topics(["multi-partition-topic"], num_partitions=3)
delete_topics(topic_names)

批量删除 Topic;Topic 不存在时由UnknownTopicOrPartitionError兜底静默忽略(kafka_utils.py)。

kafka_utils.delete_topics(["topic1", "topic2"])

生产数据

send_messages(topic, messages)

向指定 Topic 发送消息。

  • topic(str):目标 Topic 名;
  • messagesList[tuple]):(key, value)元组列表。

实现上逐条producer.send(),等待每条future.get(timeout=10)后再统一flush()(kafka_utils.py),保证返回时消息已写入 broker。

kafka_utils.send_messages("test-topic", [ ("user1", "login"), ("user2", "logout"), ("user1", "purchase"), ])

读取数据

get_all_records(spark, topic, key_deserializer="STRING", value_deserializer="STRING")

用 Spark 批量读取 Topic 的全部记录。

  • spark:SparkSession 实例;
  • topic(str):Topic 名;
  • key_deserializer(str):key 反序列化类型,默认"STRING"
  • value_deserializer(str):value 反序列化类型,默认"STRING"

实现内部使用spark.read.format("kafka"),设置startingOffsets=earliestendingOffsets=latest,再通过CAST(key AS {deserializer})做类型转换,最后sorted()排序返回(key, value)元组列表(kafka_utils.py)。排序特性让断言不受分区与消费顺序影响。

records = kafka_utils.get_all_records(self.spark, "test-topic") assert records == [("key1", "value1"), ("key2", "value2")]

测试辅助工具

assert_eventually(result_func, expected, timeout=60, interval=1.0)

轮询断言:在超时时间内反复执行result_func(),直到结果与expected相等。

  • result_func(Callable):返回当前结果的函数;
  • expected:期望结果;
  • timeout(int):最大等待秒数,默认 60;
  • interval(float):轮询间隔秒数,默认 1.0;
  • 抛出:超时后抛出AssertionError,错误信息中包含期望值与最后一次实际值(kafka_utils.py)。

它专门服务于最终一致性场景——流式查询的写端与读端之间天然存在延迟,不能像批量测试那样立即断言。

kafka_utils.assert_eventually( lambda: kafka_utils.get_all_records(self.spark, "sink-topic"), [("key1", "processed-value1")], timeout=30 )
wait_for_query_alive(query, timeout=60, interval=1.0)

等待流式查询进入活跃状态。

  • queryStreamingQuery实例;
  • timeout(int):最大等待秒数,默认 60;
  • interval(float):轮询间隔秒数,默认 1.0;
  • 抛出:超时抛出AssertionError;若查询出现异常,会直接抛出该异常

实现上循环检查query.statusisDataAvailableisTriggerActive字段,任一为真即认为查询活跃(kafka_utils.py)。因为query.exception()在非空时会立即抛出,它同时扮演了“查询失败快速暴露”的角色。

query = df.writeStream.format("memory").start() kafka_utils.wait_for_query_alive(query, timeout=30)

属性

  • broker:Kafka bootstrap server 地址,例如localhost:9093,在setup()后可用,直接用于 Spark 数据源的kafka.bootstrap.servers选项;
  • producer:底层KafkaProducer实例,供高级用法(自定义分区器、事务等)使用,访问前同样会校验初始化状态;
  • admin_client:底层KafkaAdminClient实例,可执行更细粒度的元数据管理。

常见测试模式

模式一:批量读写(Batch Read/Write)

验证 Spark DataFrame 写入 Kafka 后再读回的闭环:

def test_kafka_batch(self): topic = "test-topic" self.kafka_utils.create_topics([topic]) # Write with Spark DataFrame df = self.spark.createDataFrame([("key1", "value1")], ["key", "value"]) ( df.selectExpr("CAST(key AS BINARY)", "CAST(value AS BINARY)") .write .format("kafka") .option("kafka.bootstrap.servers", self.kafka_utils.broker) .option("topic", topic) .save() ) # Read back records = self.kafka_utils.get_all_records(self.spark, topic) assert records == [("key1", "value1")]

注意写入 Kafka 时列必须转为BINARY(对应 Kafka 记录的 key/value 字节数组),这正是 connector/kafka-0-10-sql 连接器要求的 schema。

模式二:流式查询(Streaming Queries)

Kafka 到 Kafka 的流式管道,配合 checkpoint 管理与最终一致性断言:

def test_kafka_streaming(self): import tempfile import os # Setup topics source_topic = "source" sink_topic = "sink" self.kafka_utils.create_topics([source_topic, sink_topic]) # Produce initial data self.kafka_utils.send_messages(source_topic, [("k1", "v1")]) # Start streaming query df = ( self.spark.readStream .format("kafka") .option("kafka.bootstrap.servers", self.kafka_utils.broker) .option("subscribe", source_topic) .option("startingOffsets", "earliest") .load() ) checkpoint_dir = os.path.join(tempfile.mkdtemp(), "checkpoint") query = ( df.writeStream .format("kafka") .option("kafka.bootstrap.servers", self.kafka_utils.broker) .option("topic", sink_topic) .option("checkpointLocation", checkpoint_dir) .start() ) try: self.kafka_utils.wait_for_query_alive(query) self.kafka_utils.assert_eventually( lambda: self.kafka_utils.get_all_records(self.spark, sink_topic), [("k1", "v1")] ) finally: query.stop()

这段模式与仓库真实测试 test_streaming_kafka_rtm.py 的结构一致:真实测试使用outputMode("update").trigger(realTime="30 seconds")触发流式处理,用wait_for_query_alive等待查询就绪,再用assert_eventually轮询 sink 端结果。每条测试都会用uuid生成独立的source-*/sink-*Topic 以避免用例间数据串扰(见 test_streaming_kafka_rtm.py)。

模式三:有状态聚合(Stateful Aggregations)

验证流式聚合结果:

def test_kafka_aggregation(self): # Send data for aggregation self.kafka_utils.send_messages("source", [ ("user1", "1"), ("user2", "1"), ("user1", "1"), ]) # Aggregate by key df = ( self.spark.readStream .format("kafka") .option("kafka.bootstrap.servers", self.kafka_utils.broker) .option("subscribe", "source") .load() .groupBy(col("key")) .count() .selectExpr("CAST(key AS BINARY)", "CAST(count AS STRING) AS value") ) query = df.writeStream.format("kafka") # ... start query # Verify aggregated results self.kafka_utils.assert_eventually( lambda: self.kafka_utils.get_all_records(self.spark, "sink"), [("user1", "2"), ("user2", "1")] )

有状态聚合天然是最终一致的(watermark、触发间隔都会影响何时产出聚合结果),此时assert_eventually的轮询能力尤为关键。

模式四:多 Topic 写入

通过 DataFrame 中的topic列按数据内容路由到不同 Topic:

def test_multiple_topics(self): topic1, topic2 = "topic1", "topic2" self.kafka_utils.create_topics([topic1, topic2]) # Write with topic column df = self.spark.createDataFrame([ (topic1, "key1", "value1"), (topic2, "key2", "value2"), ], ["topic", "key", "value"]) ( df.selectExpr("topic", "CAST(key AS BINARY)", "CAST(value AS BINARY)") .write .format("kafka") .option("kafka.bootstrap.servers", self.kafka_utils.broker) .save() ) # Verify data in each topic assert self.kafka_utils.get_all_records(self.spark, topic1) == [("key1", "value1")] assert self.kafka_utils.get_all_records(self.spark, topic2) == [("key2", "value2")]

当写入 schema 含topic列且未显式指定topic选项时,Kafka 连接器按行路由目标 Topic,这是 Kafka sink 的内置能力。

运行测试

运行全部 Kafka 测试

cd $SPARK_HOME/python python -m pytest pyspark/sql/tests/streaming/test_streaming_kafka_rtm.py -v

运行单个用例

python -m pytest pyspark/sql/tests/streaming/test_streaming_kafka_rtm.py::StreamingKafkaTests::test_streaming_stateless -v

使用 unittest 运行

cd $SPARK_HOME/python python -m unittest pyspark.sql.tests.streaming.test_streaming_kafka_rtm

需要说明的依赖前提:测试类StreamingKafkaTests上有三个@unittest.skipIf条件(test_streaming_kafka_rtm.py),任一不满足即整类跳过:

  1. 未安装kafka包(报No module named 'kafka');
  2. 未安装testcontainers包(报No module named 'testcontainers');
  3. Docker daemon 不可用(报Docker is not available)。

因此建议先完成前文的环境准备,再运行测试,才能看到真实的容器启动与流式断言过程。

故障排查

Docker 未运行

报错:

Cannot connect to the Docker daemon at unix:///var/run/docker.sock

解决:启动 Docker Desktop 或 Docker daemon,确认docker ps可正常输出。

容器启动超时

报错:

Kafka container failed to start within timeout

解决:

  1. 在测试代码中增大超时(testcontainers 侧可配置容器等待超时);
  2. 检查 Docker 资源分配(CPU/内存),镜像拉取与 broker 启动都较吃资源;
  3. 查看容器日志定位问题:docker logs <container-id>

端口冲突

报错:

Port 9093 already in use

解决:testcontainers 会自动分配随机映射端口,不要在测试里手工绑定固定端口;若本地残留旧容器,清理后重试。这也是文档强调“broker 地址通过kafka_utils.broker动态获取”的原因。

依赖缺失

报错:

ImportError: testcontainers is required for Kafka tests

解决

pip install testcontainers[kafka] kafka-python-ng

Kafka 连接器 JAR 缺失

报错:运行测试时抛出Kafka SQL connector JAR was not found

解决:先构建 Spark 产物再运行测试:

build/mvn package # 或 build/sbt Test/package

深入理解:底层实现要点

1. broker 地址的动态获取

testcontainers 每次启动容器都映射不同的宿主机端口,get_bootstrap_server()返回的地址形如localhost:9093。测试中始终通过self.kafka_utils.broker读取,而不是硬编码,这正是“端口冲突”问题能被系统性规避的原因。

2. 生产者序列化策略

setup()中构建的KafkaProducer使用str(k).encode("utf-8")对 key 和 value 做序列化(kafka_utils.py)。因此send_messages接受任意可str()的对象,get_all_records读回的字符串与str()结果一致,例如send_messages(topic, [(i, i) for i in range(10)])读回的是("0", "0") ... ("9", "9")(与真实测试用例 test_streaming_kafka_rtm.py 的expected = sorted((str(i), str(i)) for i in range(10))对应)。

3. 断言策略:排序 + 最终一致

get_all_records对结果排序返回,使断言与 Kafka 分区顺序解耦;assert_eventually的轮询机制则吸收流式处理端到端延迟。两者组合是 Kafka 流式测试中“稳定断言”的黄金搭配。

4. 与 PySpark 测试体系的契合

KafkaUtils被设计为ReusedSQLTestCase的类级夹具:setUpClasssetup()tearDownClassteardown(),单个测试方法内再用setUp/tearDown创建与清理每次用例的独立 Topic。仓库的StreamingKafkaTestsMixin(test_streaming_kafka_rtm.py)把这一整套模式沉淀成了可复用的 mixin,新测试类只需继承它即可获得完整的 Kafka 测试能力。

结语

通过KafkaUtils与 testcontainers,PySpark 开发者可以把“真实 Kafka 集群 + 端到端流式断言”变成一条pytest命令即可完成的本地位测试流程。本文覆盖的环境准备、API 全览、四种测试模式与故障排查,直接对应仓库中的 KAFKA_TESTING.md、kafka_utils.py 与 test_streaming_kafka_rtm.py,你可以直接以此为模板,为你的 Spark Kafka 应用编写同样可靠的集成测试。

  • 大数据
  • 数据分析
  • 批处理
  • 流处理
  • 机器学习
  • 图计算

【免费下载链接】spark

Apache Spark - A unified analytics engine for large-scale data processing

项目地址:https://gitcode.com/gh_mirrors/sp/spark
点击查看免费下载

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询