- 大数据
- 数据分析
- 批处理
- 流处理
- 机器学习
- 图计算
【免费下载链接】spark
Apache Spark - A unified analytics engine for large-scale data processing
导读
本指南面向需要在本地对 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 devsetup()源码中(kafka_utils.py)对依赖做了显式校验:缺少testcontainers.kafka或kafka(KafkaProducer/KafkaAdminClient)都会抛出带安装提示的ImportError。测试端对应的跳检逻辑在 python/pyspark/testing/utils.py,通过have_package探测kafka与testcontainers,缺失时由@unittest.skipIf跳过测试类。
3. Spark 构建产物
运行 Kafka 测试还需要 Spark 的 Kafka SQL 连接器 JAR。测试基类StreamingKafkaTestsMixin在创建 SparkSession 之前会通过search_jar在connector/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):
- 幂等检查:若
initialized为 True 直接返回; - 导入
testcontainers.kafka.KafkaContainer,缺失则抛ImportError; - 导入
kafka.KafkaProducer与kafka.admin.KafkaAdminClient,缺失则抛ImportError; - 启动容器,通过
get_bootstrap_server()获得broker地址(形如localhost:9093,端口由 testcontainers 随机映射); - 用 10 秒超时初始化
KafkaAdminClient和KafkaProducer(producer 的 key/value 序列化器会把非 None 值str()后 UTF-8 编码); - 任意一步异常都会先调用
teardown()清理再重新抛出RuntimeError。
可能抛出的异常:
ImportError:依赖未安装;RuntimeError:容器启动失败。
注意:首次运行时 Docker 需要拉取镜像,容器启动可能耗时 10~30 秒。
teardown()
关闭 admin client、关闭 producer(5 秒超时)、停止容器并复位所有状态字段。实现上每一步都做了异常兜底与字段置空(kafka_utils.py),因此可安全多次调用,适合放在finally或tearDownClass中。
内部防护:_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_names(
List[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 名;
- messages(
List[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=earliest、endingOffsets=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)
等待流式查询进入活跃状态。
- query:
StreamingQuery实例; - timeout(int):最大等待秒数,默认 60;
- interval(float):轮询间隔秒数,默认 1.0;
- 抛出:超时抛出
AssertionError;若查询出现异常,会直接抛出该异常。
实现上循环检查query.status的isDataAvailable或isTriggerActive字段,任一为真即认为查询活跃(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),任一不满足即整类跳过:
- 未安装
kafka包(报No module named 'kafka'); - 未安装
testcontainers包(报No module named 'testcontainers'); - 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解决:
- 在测试代码中增大超时(testcontainers 侧可配置容器等待超时);
- 检查 Docker 资源分配(CPU/内存),镜像拉取与 broker 启动都较吃资源;
- 查看容器日志定位问题:
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-ngKafka 连接器 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的类级夹具:setUpClass中setup()、tearDownClass中teardown(),单个测试方法内再用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
相关推荐
JUnit4与Apache Kafka Streams集成:流处理测试
JUnit4与Apache Kafka Streams集成:流处理测试 引言:流处理测试的痛点与解决方案 你是否在开发Apache Kafka Streams应
测试开发工具Apache Kafka 系统级测试指南:基于 ducktape 在 Docker、本地虚拟机与 EC2 上运行集成与性能测试
Apache Kafka 系统级测试指南:基于 ducktape 在 Docker、本地虚拟机与 EC2 上运行集成与性能测试 Apache Kafka 仓库中
消息队列流处理数据集成存储快速拉取 App Store IPA 包:ipatool 完整下载教程
快速拉取 App Store IPA 包:ipatool 完整下载教程 需要从 App Store 批量拉取 IPA 包做测试、备份,又不想在图形界面里反复点按
CLI开发工具
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考