资讯详情

资讯详情

PySpark Kafka 集成测试指南:用 KafkaUtils 与 Docker 容器在本地测试 Kafka 流处理

大数据数据分析批处理流处理机器学习图计算【免费下载链接】sparkApache 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()关闭容器并清理客户端资源测试结束不留残留进程贴近生产测试的是真实 brokerSpark 读到的 key/value 是真实的二进制数据行为与生产环境一致。从源码结构看这一工具的设计意图是与 Python unittest 体系尤其是ReusedSQLTestCase深度配合作为整个测试类的类级夹具使用。环境准备1. DockerKafkaUtils依赖 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或kafkaKafkaProducer/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 不存在会提示先构建 Sparkbuild/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_version7.4.0)创建一个KafkaUtils实例。kafka_version使用的 Confluent Kafka 镜像版本默认7.4.0。源码在setup()中将其拼为confluentinc/cp-kafka:{kafka_version}传给KafkaContainerkafka_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和KafkaProducerproducer 的 key/value 序列化器会把非 None 值str()后 UTF-8 编码任意一步异常都会先调用teardown()清理再重新抛出RuntimeError。可能抛出的异常ImportError依赖未安装RuntimeError容器启动失败。注意首次运行时 Docker 需要拉取镜像容器启动可能耗时 10~30 秒。teardown()关闭 admin client、关闭 producer5 秒超时、停止容器并复位所有状态字段。实现上每一步都做了异常兜底与字段置空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_partitions1, replication_factor1)批量创建 Topic。topic_namesList[str]要创建的 Topic 名列表num_partitionsint每个 Topic 的分区数默认 1replication_factorint副本因子默认 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_partitions3)delete_topics(topic_names)批量删除 TopicTopic 不存在时由UnknownTopicOrPartitionError兜底静默忽略kafka_utils.py。kafka_utils.delete_topics([topic1, topic2])生产数据send_messages(topic, messages)向指定 Topic 发送消息。topicstr目标 Topic 名messagesList[tuple](key, value)元组列表。实现上逐条producer.send()等待每条future.get(timeout10)后再统一flush()kafka_utils.py保证返回时消息已写入 broker。kafka_utils.send_messages(test-topic, [ (user1, login), (user2, logout), (user1, purchase), ])读取数据get_all_records(spark, topic, key_deserializerSTRING, value_deserializerSTRING)用 Spark 批量读取 Topic 的全部记录。sparkSparkSession 实例topicstrTopic 名key_deserializerstrkey 反序列化类型默认STRINGvalue_deserializerstrvalue 反序列化类型默认STRING。实现内部使用spark.read.format(kafka)设置startingOffsetsearliest、endingOffsetslatest再通过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, timeout60, interval1.0)轮询断言在超时时间内反复执行result_func()直到结果与expected相等。result_funcCallable返回当前结果的函数expected期望结果timeoutint最大等待秒数默认 60intervalfloat轮询间隔秒数默认 1.0抛出超时后抛出AssertionError错误信息中包含期望值与最后一次实际值kafka_utils.py。它专门服务于最终一致性场景——流式查询的写端与读端之间天然存在延迟不能像批量测试那样立即断言。kafka_utils.assert_eventually( lambda: kafka_utils.get_all_records(self.spark, sink-topic), [(key1, processed-value1)], timeout30 )wait_for_query_alive(query, timeout60, interval1.0)等待流式查询进入活跃状态。queryStreamingQuery实例timeoutint最大等待秒数默认 60intervalfloat轮询间隔秒数默认 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, timeout30)属性brokerKafka 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 QueriesKafka 到 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(realTime30 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列按数据内容路由到不同 Topicdef 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 testcontainersDocker 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。仓库的StreamingKafkaTestsMixintest_streaming_kafka_rtm.py把这一整套模式沉淀成了可复用的 mixin新测试类只需继承它即可获得完整的 Kafka 测试能力。结语通过KafkaUtils与 testcontainersPySpark 开发者可以把“真实 Kafka 集群 端到端流式断言”变成一条pytest命令即可完成的本地位测试流程。本文覆盖的环境准备、API 全览、四种测试模式与故障排查直接对应仓库中的 KAFKA_TESTING.md、kafka_utils.py 与 test_streaming_kafka_rtm.py你可以直接以此为模板为你的 Spark Kafka 应用编写同样可靠的集成测试。赞分享大数据数据分析批处理流处理机器学习图计算【免费下载链接】sparkApache Spark - A unified analytics engine for large-scale data processing项目地址https://gitcode.com/gh_mirrors/sp/spark点击查看免费下载相关推荐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),仅供参考
觉得有用,分享给同行:

为您的企业打造数字门面

稳重轻奢商务风格,端正雅致视觉,长效耐看不易过时。

立即咨询 →