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

如何用Scala合并多同源Flink Job输出并聚合为指定统一格式?

解决方案

1. 输出主题规划

给每个Flink Job分配独立的Kafka输出主题(比如job1-output、job2-output...jobn-output),聚合Job一次性订阅所有这些主题,避免不同Job的输出互相干扰。

2. 定义数据实体类

用Scala case class统一输入输出的结构:

// 单个Flink Job的输出结构
case class JobOutput(pk: String, variable1: String, variable2: Boolean)

// 最终聚合后的输出结构
case class AggregatedOutput(pk: String, variable1: List[String], variable2: List[Boolean])

3.1 消费所有Job的输出主题

使用Flink Kafka Consumer订阅多个输出主题,并反序列化数据:

import org.apache.flink.api.common.serialization.SimpleStringSchema
import org.apache.flink.streaming.api.scala._
import org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumer
import com.fasterxml.jackson.databind.ObjectMapper
import com.fasterxml.jackson.module.scala.DefaultScalaModule

// 初始化Flink流处理环境
val env = StreamExecutionEnvironment.getExecutionEnvironment

// Kafka连接配置
val kafkaProps = new java.util.Properties()
kafkaProps.setProperty("bootstrap.servers", "your-kafka-broker:9092")
kafkaProps.setProperty("group.id", "aggregation-job-group")

// 订阅所有Job的输出主题,可根据实际数量调整
val outputTopics = List("job1-output", "job2-output", "job3-output")
val kafkaSource = new FlinkKafkaConsumer[String](outputTopics, new SimpleStringSchema(), kafkaProps)

// 将Kafka中的JSON字符串反序列化为JobOutput对象
val jobOutputStream = env.addSource(kafkaSource)
  .map { jsonStr =>
    val mapper = new ObjectMapper().registerModule(DefaultScalaModule)
    mapper.readValue(jsonStr, classOf[JobOutput])
  }

3.2 按pk聚合数据

使用KeyedProcessFunction维护状态,收集每个pk对应的variable1和variable2列表:

import org.apache.flink.api.common.state.{ListState, ListStateDescriptor}
import org.apache.flink.streaming.api.functions.KeyedProcessFunction
import org.apache.flink.util.Collector

class AggregationFunction extends KeyedProcessFunction[String, JobOutput, AggregatedOutput] {
  // 存储每个pk对应的variable1列表
  private lazy val variable1State: ListState[String] = getRuntimeContext.getListState(
    new ListStateDescriptor[String]("variable1-state", classOf[String])
  )
  // 存储每个pk对应的variable2列表
  private lazy val variable2State: ListState[Boolean] = getRuntimeContext.getListState(
    new ListStateDescriptor[Boolean]("variable2-state", classOf[Boolean])
  )

  override def processElement(
    value: JobOutput,
    ctx: KeyedProcessFunction[String, JobOutput, AggregatedOutput]#Context,
    out: Collector[AggregatedOutput]
  ): Unit = {
    // 将当前记录的字段加入对应状态
    variable1State.add(value.variable1)
    variable2State.add(value.variable2)

    // 注册10秒后的处理时间定时器,触发聚合结果输出(可根据业务调整触发逻辑)
    ctx.timerService().registerProcessingTimeTimer(ctx.timerService().currentProcessingTime() + 10000)
  }

  override def onTimer(
    timestamp: Long,
    ctx: KeyedProcessFunction[String, JobOutput, AggregatedOutput]#OnTimerContext,
    out: Collector[AggregatedOutput]
  ): Unit = {
    // 从状态中提取所有元素并转为列表,可根据需求添加去重逻辑
    val var1List = variable1State.get().iterator().asScala.toList.distinct
    val var2List = variable2State.get().iterator().asScala.toList.distinct

    // 输出聚合结果
    out.collect(AggregatedOutput(ctx.getCurrentKey, var1List, var2List))

    // 清空状态,避免重复输出(可选,根据业务是否需要保留状态调整)
    variable1State.clear()
    variable2State.clear()
  }
}

// 按pk分组并应用聚合逻辑
val aggregatedStream = jobOutputStream
  .keyBy(_.pk)
  .process(new AggregationFunction)

3.3 输出JSON格式结果

将聚合后的对象序列化为JSON字符串,输出到目标存储(示例为标准输出,可替换为Kafka Sink、文件等):

// 将聚合结果序列化为JSON字符串
val jsonStream = aggregatedStream
  .map { aggOutput =>
    val mapper = new ObjectMapper().registerModule(DefaultScalaModule)
    mapper.writeValueAsString(aggOutput)
  }

// 输出到标准输出(实际场景替换为业务需要的Sink)
jsonStream.print()

// 执行聚合Job
env.execute("Flink Job Output Aggregation")

4. 关键注意事项

  • 状态与触发逻辑:如果需要等待所有Job的输出再聚合,可改用事件时间结合水位线判断,替代处理时间定时器。
  • 去重处理:如果存在重复输出的场景,可在聚合时添加去重逻辑(如示例中的distinct)。
  • 容错保障:开启Flink的Checkpoint机制,确保故障时状态可恢复,避免数据丢失。
  • 扩展性:新增Flink Job时,只需将新的输出主题加入订阅列表,无需修改聚合Job核心逻辑。

内容的提问来源于stack exchange,提问作者Hardik Doshi

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 22:55:23