You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spark Streaming订阅多Kafka主题后如何写入对应目标主题?

解决方案

要实现将数据写入对应名称的Kafka主题,核心是利用Spark Kafka Sink的动态主题写入特性——通过DataFrame中的topic字段指定每条数据的输出目标,而非固定单一主题。具体实现如下:

关键修改点

  1. 保留并转换原始主题名称:从输入的topic字段生成对应的输出主题(将前缀ingestion_替换为processed_)
  2. 移除固定主题配置:不再通过.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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.25 07:36:17