Spark Streaming订阅多Kafka主题后如何写入对应目标主题?
解决方案
要实现将数据写入对应名称的Kafka主题,核心是利用Spark Kafka Sink的动态主题写入特性——通过DataFrame中的topic字段指定每条数据的输出目标,而非固定单一主题。具体实现如下:
关键修改点
- 保留并转换原始主题名称:从输入的
topic字段生成对应的输出主题(将前缀ingestion_替换为processed_) - 移除固定主题配置:不再通过
.option("topic", ...)指定单一输出主题,让Kafka Sink自动使用DataFrame中的topic列作为目标
修改后的完整代码
(spark.readStream.format("kafka") .option("kafka.bootstrap.servers", bootstrap_server) .option("subscribePattern", "ingestion_src_api.*") .option("startingOffsets", "latest") .load() # 保留原始topic字段,解析Kafka消息value .select(col("topic").cast("string"), from_json(col("value").cast("string"), schema).alias("value")) # 生成目标topic列,同时保留转换后的value .select( # 将原始主题前缀替换,得到对应输出主题 regexp_replace(col("topic"), "^ingestion_", "processed_").alias("topic"), # 原有value转换逻辑保持不变 to_json(struct( expr("value.active_id as active_id"), expr("value.size as timeframe"), expr("cast(value.at / 1000000000 as timestamp) as executed_at"), expr("FROM_UNIXTIME(value.from) as candle_from"), expr("FROM_UNIXTIME(value.to) as candle_to"), expr("value.id as period"), "value.open", "value.close", "value.min", "value.max", "value.ask", "value.bid", "value.volume" )).alias("value") ) # 写入Kafka,移除固定topic配置 .writeStream.format("kafka") .option("kafka.bootstrap.servers", bootstrap_server) .option("checkpointLocation", "./checkpoint/") .start() )
原理说明
- Spark Kafka Sink会自动识别DataFrame中的
topic(字符串类型)和value(字符串/二进制类型)列,将每条数据写入topic列指定的Kafka主题 - 通过
regexp_replace直接替换主题前缀,确保输入主题与输出主题一一对应(如ingestion_src_api_iq_BTCUSD_1_json→processed_src_api_iq_BTCUSD_1_json) - 检查点目录会自动维护不同主题的消费偏移量,无需额外配置
内容的提问来源于stack exchange,提问作者Renan Nogueira
相关产品推荐
相关产品推荐

