如何流式读取Kafka单Topic按Key过滤后写入多个HDFS路径?
问题根因
现有代码存在4个核心问题,导致无法正常运行:
- 重复创建Kafka消费链路:针对同一个Topic1启动了2个独立的
readStream,相当于启动2个独立消费者组重复拉取全量Topic数据,会造成资源浪费,数据量较大时会直接触发消费延迟、Executor OOM导致作业失败。 - Key类型匹配逻辑错误:Spark读取Kafka得到的
key字段默认是二进制(BinaryType)类型,直接和字符串"Value1"/"Value2"做等值判断永远返回false,过滤后无数据写入目标路径。 - 缺少作业阻塞逻辑:所有流查询调用
start()后没有加终止等待逻辑,Spark作业启动后会直接退出,无法持续消费数据。 - 存在冗余计算:多次重复对
value字段做字符串类型转换,浪费计算资源。
正确实现方案
核心思路为单次读取Kafka Topic,按Key规则分流后写入多HDFS路径,避免重复消费,修复代码如下:
import org.apache.spark.sql.streaming.Trigger import org.apache.spark.sql.functions.col // 仅需读取1次Kafka Topic,从源头避免重复消费 val kafkaStream = spark .readStream .format("kafka") .option("kafka.bootstrap.servers", configManager.getString("Kafka.Server")) .option("subscribe", "Topic1") .option("startingOffsets", "latest") .option("failOnDataLoss", "false") .load() // 提前统一做字段类型转换,消除后续冗余计算 .select( // 核心修复:将二进制类型的key转为字符串,保证过滤逻辑生效 col("key").cast("string").as("key"), col("value").cast("string").as("value") ) // 分流处理Key=Value1的数据,写入对应HDFS路径 val value1Query = kafkaStream .filter(col("key") === "Value1") .select(functions.from_json(col("value"), Value1Schema.schemaExecution).as("value")) .select("value.*") .writeStream .format("orc") .option("metastoreUri", configManager.getString("spark.datasource.hive.warehouse.metastoreUri")) // 不同写入流的checkpoint路径必须完全独立,禁止共用 .option("checkpointLocation", "/tmp/teststreaming/execution/checkpoint2005") .option("path", "/tmp/test/value1") .trigger(Trigger.ProcessingTime("5 Seconds")) .partitionBy("jobid") .start() // 分流处理Key=Value2的数据,写入对应HDFS路径 val value2Query = kafkaStream .filter(col("key") === "Value2") .select(functions.from_json(col("value"), Value2Schema.schemaJobParameters).as("value")) .select("value.*") .writeStream .format("orc") .option("metastoreUri", configManager.getString("spark.datasource.hive.warehouse.metastoreUri")) .option("checkpointLocation", "/tmp/teststreaming/jobparameters/checkpoint2006") .option("path", "/tmp/test/value2") .trigger(Trigger.ProcessingTime("5 Seconds")) .partitionBy("jobid") .start() // 阻塞等待所有流作业运行,避免启动后直接退出 spark.streams.awaitAnyTermination()
运行注意事项
- 首次部署前必须清空对应checkpoint目录,旧的状态数据会导致作业从历史offset恢复,出现数据不符合预期的问题。
- 如果Kafka Key存在null值,可以在filter逻辑前增加
.where(col("key").isNotNull),避免空值导致的计算异常。 - 后续如果需要新增更多Key的分流规则,直接基于同一个
kafkaStream做过滤新增写入流即可,不需要重复创建Kafka读链路。
内容的提问来源于stack exchange,提问作者dataeng
相关产品推荐
相关产品推荐

