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

Flink Kafka Source算子并行度调整效果咨询

Flink作业并行度与输出速率不匹配问题分析

作业DAG结构

  • 包含4个输出Sink
  • 通过Kafka Source读取单个Topic,Topic的每个Partition已按Key分区
  • 执行链路:
    1. Source -> KeyedProcessFunction -> Map -> Sink-1
    2. Source -> KeyedProcessFunction -> Map -> FlatMap -> Sink-2
    3. Source -> KeyedProcessFunction -> Map -> FlatMap -> Sink-3
    4. Source -> KeyedProcessFunction -> Map -> Sink-4

问题描述

将Kafka Source算子的并行度设置为与输入Kafka Partition数量匹配后,并行度会产生差异吗?
备注:观察到Kafka Source的数据消费速率有所提升,但输出速率与输入速率不匹配。

作业代码

class JobFlowGraph(val config: Config) {

  private final val env = StreamExecutionEnvironment.getExecutionEnvironment

  private final val streamingRecordType: TypeInformation[StreamingRecord] = createTypeInformation[StreamingRecord]

  private final val somaLogRecordType: AvroTypeInfo[SomaLogRecordAvro] = new AvroTypeInfo[SomaLogRecordAvro](classOf[SomaLogRecordAvro])

  private def kafkaSource(uid: String) = env.fromSource(getKafkaSource,
    WatermarkStrategy
      .forBoundedOutOfOrderness[StreamingRecord](Duration.ofMinutes(1))
      .withTimestampAssigner(new SerializableTimestampAssigner[StreamingRecord] {
        override def extractTimestamp(element: StreamingRecord, recordTimestamp: Long): Long = element.eventTs
      })
      .withIdleness(Duration.ofMinutes(10))
  , "Kafka Source")(streamingRecordType)
    .uid(uid+config.get(KAFKA_SOURCE_UID))

  private def getKafkaAuthProps: Properties = {
    val commonKafkaProps = new Properties()
    commonKafkaProps.put(CommonClientConfigs.SECURITY_PROTOCOL_CONFIG, SecurityProtocol.SASL_SSL.name)
    commonKafkaProps.put(SaslConfigs.SASL_MECHANISM, "AWS_MSK_IAM")
    commonKafkaProps.put(SaslConfigs.SASL_JAAS_CONFIG, "software.amazon.msk.auth.iam.IAMLoginModule required;")
    commonKafkaProps.put(SaslConfigs.SASL_CLIENT_CALLBACK_HANDLER_CLASS, "software.amazon.msk.auth.iam.IAMClientCallbackHandler")
    commonKafkaProps
  }

  private def getKafkaConsumerProps: Properties = {
    val consumerKafkaProps = new Properties()
    consumerKafkaProps.putAll(getKafkaAuthProps)
    consumerKafkaProps.put(ConsumerConfig.GROUP_ID_CONFIG, config.get(KAFKA_CONSUMER_GROUP_ID).concat(config.get(REGION)))
    consumerKafkaProps.put(ConsumerConfig.CLIENT_ID_CONFIG, config.get(KAFKA_CONSUMER_CLIENT_ID).concat(config.get(REGION)))
    consumerKafkaProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, config.get(OFFSET_ON_AUTO_RESET))
    consumerKafkaProps.put(ConsumerConfig.ISOLATION_LEVEL_CONFIG, IsolationLevel.READ_COMMITTED.name().toLowerCase(Locale.ROOT))
    consumerKafkaProps.put(ConsumerConfig.REQUEST_TIMEOUT_MS_CONFIG, config.get(KAFKA_REQUEST_TIMEOUT_MS))
    consumerKafkaProps
  }

  private def getKafkaProducerProps: Properties = {
    val kafkaDataVizDPIProducerProps = new Properties()
    kafkaDataVizDPIProducerProps.putAll(getKafkaAuthProps)
    kafkaDataVizDPIProducerProps.put("enable.idempotence", "true")
    kafkaDataVizDPIProducerProps.put("transaction.timeout.ms", "900000") // should be same as in MSK cluster transaction timeout
    kafkaDataVizDPIProducerProps
  }


  private def getKafkaSource: KafkaSource[StreamingRecord]  = {
    KafkaSource.builder[StreamingRecord]
      .setBootstrapServers(config.get(INPUT_KAFKA_HOST))
      .setTopics(config.get(INPUT_TOPIC))
      .setProperties(getKafkaConsumerProps)
      .setProperty("partition.discovery.interval.ms", "100000") // discover new partitions per 100 seconds
      .setStartingOffsets(OffsetsInitializer.committedOffsets(OffsetResetStrategy.EARLIEST))
      .setDeserializer(new RecordKafkaDeserializationSchema())
      .build()
  }

  private def getKafkaDPISink: KafkaSink[DPIKafkaOutputRecord] = {
    val schemaRegistryConfig: mutable.HashMap[String, Any] = mutable.HashMap()
    schemaRegistryConfig.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, classOf[org.apache.kafka.common.serialization.StringSerializer])
    schemaRegistryConfig.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG,
      classOf[io.confluent.kafka.serializers.KafkaAvroSerializer])
    schemaRegistryConfig.put("schema.registry.url", config.get(SCHEMA_REGISTRY_URL))
    schemaRegistryConfig.put("use.schema.id", config.get(SCHEMA_REGISTRY_DPI_SUBJECT_VERSION))
    schemaRegistryConfig.put("auto.register.schemas", false)
    schemaRegistryConfig.put("use.latest.version", false)

    KafkaSink.builder[DPIKafkaOutputRecord]
      .setBootstrapServers(config.get(OUTPUT_KAFKA_HOST))
      .setKafkaProducerConfig(getKafkaProducerProps)
      .setDeliveryGuarantee(DeliveryGuarantee.EXACTLY_ONCE)
      .setTransactionalIdPrefix(config.get(KAFKA_PRODUCER_DV_DPI_TRANSACTION_ID))
      .setRecordSerializer(new RecordKafkaSerializationSchema[DPIKafkaOutputRecord](config.get(SCHEMA_REGISTRY_DPI_SUBJECT), DVDemandPartnerInteractionAvro.getClassSchema, config.get(SCHEMA_REGISTRY_URL), schemaRegistryConfig))
      .build()
  }

  private def getKafkaAuctionSink: KafkaSink[AuctionKafkaOutputRecord] = {
    val schemaRegistryConfig: mutable.HashMap[String, Any] = mutable.HashMap()
    schemaRegistryConfig.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, classOf[org.apache.kafka.common.serialization.StringSerializer])
    schemaRegistryConfig.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG,
      classOf[io.confluent.kafka.serializers.KafkaAvroSerializer])
    schemaRegistryConfig.put("schema.registry.url", config.get(SCHEMA_REGISTRY_URL))
    schemaRegistryConfig.put("use.schema.id", config.get(SCHEMA_REGISTRY_AUCTION_SUBJECT_VERSION))
    schemaRegistryConfig.put("auto.register.schemas", false)
    schemaRegistryConfig.put("use.latest.version", false)

    KafkaSink.builder[AuctionKafkaOutputRecord]
      .setBootstrapServers(config.get(OUTPUT_KAFKA_HOST))
      .setKafkaProducerConfig(getKafkaProducerProps)
      .setDeliveryGuarantee(DeliveryGuarantee.EXACTLY_ONCE)
      .setTransactionalIdPrefix(config.get(KAFKA_PRODUCER_DV_AUCTION_TRANSACTION_ID))
      .setRecordSerializer(new RecordKafkaSerializationSchema[AuctionKafkaOutputRecord](config.get(SCHEMA_REGISTRY_AUCTION_SUBJECT), DVAuctionAvro.getClassSchema, config.get(SCHEMA_REGISTRY_URL), schemaRegistryConfig))
      .build()
  }

  private def getKafkaSessionizerSink: KafkaSink[SessionizerKafkaOutputRecord] = {
    val schemaRegistryConfig: mutable.HashMap[String, Any] = mutable.HashMap()
    schemaRegistryConfig.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, classOf[org.apache.kafka.common.serialization.StringSerializer])
    schemaRegistryConfig.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG,
      classOf[io.confluent.kafka.serializers.KafkaAvroSerializer])
    schemaRegistryConfig.put("schema.registry.url", config.get(SCHEMA_REGISTRY_URL))
    schemaRegistryConfig.put("use.schema.id", config.get(SCHEMA_REGISTRY_SOMALOGRECORD_SUBJECT_VERSION))
    schemaRegistryConfig.put("auto.register.schemas", false)
    schemaRegistryConfig.put("use.latest.version", false)

    KafkaSink.builder[SessionizerKafkaOutputRecord]
      .setBootstrapServers(config.get(OUTPUT_KAFKA_HOST))
      .setKafkaProducerConfig(getKafkaProducerProps)
      .setDeliveryGuarantee(DeliveryGuarantee.EXACTLY_ONCE)
      .setTransactionalIdPrefix(config.get(KAFKA_PRODUCER_SESSIONIZER_OUT_TRANSACTION_ID))
      .setRecordSerializer(new RecordKafkaSerializationSchema[SessionizerKafkaOutputRecord](config.get(SCHEMA_REGISTRY_SESSIONIZER_OUT_SUBJECT), SomaLogRecordAvro.getClassSchema, config.get(SCHEMA_REGISTRY_URL), schemaRegistryConfig))
      .build()
  }


  private def getKafkaSessionzierNonPIISink: KafkaSink[SessionzierNonPIIKafkaOutputRecord] = {
    val schemaRegistryConfig: mutable.HashMap[String, Any] = mutable.HashMap()
    schemaRegistryConfig.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, classOf[org.apache.kafka.common.serialization.StringSerializer])
    schemaRegistryConfig.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG,
      classOf[io.confluent.kafka.serializers.KafkaAvroSerializer])
    schemaRegistryConfig.put("schema.registry.url", config.get(SCHEMA_REGISTRY_URL))
    schemaRegistryConfig.put("use.schema.id", config.get(SCHEMA_REGISTRY_SOMALOGRECORD_SUBJECT_VERSION))
    schemaRegistryConfig.put("auto.register.schemas", false)
    schemaRegistryConfig.put("use.latest.version", false)

    KafkaSink.builder[SessionzierNonPIIKafkaOutputRecord]
      .setBootstrapServers(config.get(OUTPUT_KAFKA_HOST))
      .setKafkaProducerConfig(getKafkaProducerProps)
      .setDeliveryGuarantee(DeliveryGuarantee.EXACTLY_ONCE)
      .setTransactionalIdPrefix(config.get(KAFKA_PRODUCER_SESSIONIZER_NONPII_OUT_TRANSACTION_ID))
      .setRecordSerializer(new RecordKafkaSerializationSchema[SessionzierNonPIIKafkaOutputRecord](config.get(SCHEMA_REGISTRY_SESSIONIZER_NONPII_OUT_SUBJECT), SomaLogRecordAvro.getClassSchema, config.get(SCHEMA_REGISTRY_URL), schemaRegistryConfig))
      .build()
  }

  private def getSourceDataStream(uid: String): DataStream[StreamingRecord] = kafkaSource(uid)

  private def getSessionizedDataStream(uid: String): DataStream[SomaLogRecordAvro] = {
    getSourceDataStream(uid)
      .keyBy[String]((record: StreamingRecord) => record.sessionId)(createTypeInformation[String])
      .process[SomaLogRecordAvro](new CustomKeyedProcessFunction())(somaLogRecordType)
      .uid(uid+config.get(CUSTOM_KEY_PROCESS_FN_UID))
  }


  private def sinkToKafka(topic: String): Unit = {
    val schemaRegistryAuctionSubjectVersion: String = config.get(SCHEMA_REGISTRY_AUCTION_SUBJECT_VERSION).toString
    val schemaRegistryDPISubjectVersion: String = config.get(SCHEMA_REGISTRY_DPI_SUBJECT_VERSION).toString

    if (topic.equalsIgnoreCase(config.get(DV_DEMAND_PARTNER_INTERACTION_TOPIC))) {
      getSessionizedDataStream("dpi_task")
        .flatMap[DVDemandPartnerInteractionAvro]((rec: SomaLogRecordAvro, out: Collector[DVDemandPartnerInteractionAvro]) => {
          try {
            val records: List[DVDemandPartnerInteractionAvro] = DemandPartnerInteractionConverterTask.process(rec)
            if (records != null && records.nonEmpty) {
              records.foreach(dpiRecord => {
                if (dpiRecord == null) logger.info("null DPI record!") else {
                  dpiRecord.getMetadata.setSchemaId(schemaRegistryDPISubjectVersion)
                  out.collect(dpiRecord)
                }
              })
            }
          } catch {
            case ex: Exception => ex.printStackTrace()
          }
        })
        .map[DPIKafkaOutputRecord]((rec: DVDemandPartnerInteractionAvro)  => DPIKafkaOutputRecord(
          "", topic, rec.getMetadata.getSessionEventTs, rec)
        )
        .sinkTo(getKafkaDPISink)
        .name("dpi_kafka_sink")
    }

    if (topic.equalsIgnoreCase(config.get(DV_AUCTION_TOPIC))) {
      getSessionizedDataStream("auction_task")
        .flatMap[DVAuctionAvro]((rec: SomaLogRecordAvro, out: Collector[DVAuctionAvro]) => {
          try {
            val records: List[DVAuctionAvro] = AuctionConverterTask.process(rec)
            if (records != null && records.nonEmpty) {
              records.foreach(auctionRecord => {
                if (auctionRecord == null) logger.info("null auction record!") else {
                  auctionRecord.getMetadata.setSchemaId(schemaRegistryAuctionSubjectVersion)
                  out.collect(auctionRecord)
                }
              })
            }
          } catch {
            case ex: Exception => ex.printStackTrace()
          }
        })
        .map[AuctionKafkaOutputRecord]((rec: DVAuctionAvro) => AuctionKafkaOutputRecord(
          "", topic, rec.getMetadata.getSessionEventTs, rec)
        )
        .sinkTo(getKafkaAuctionSink)
        .name("auction_kafka_sink")
    }

    if (topic.equalsIgnoreCase(config.get(SESSIONIZER_OUT_TOPIC))) {
      getSessionizedDataStream("pii_task")
        .map[SessionizerKafkaOutputRecord]((rec: SomaLogRecordAvro) => SessionizerKafkaOutputRecord(
          "", topic, rec.getMetadata.getSessionEventTs, rec)
        )
        .sinkTo(getKafkaSessionizerSink)
        .name("sessionizer_out_kafka_sink")
    }

    if (topic.equalsIgnoreCase(config.get(SESSIONIZER_OUT_NON_PII_TOPIC))) {
      val salt = config.get(NONPIITASK_SALT_STATIC_VALUE)
      getSessionizedDataStream("nonpii_task")
        .map[SessionzierNonPIIKafkaOutputRecord]((rec: SomaLogRecordAvro) => {
          val nonPIIRecord = NonPIITask.process(rec, 4, salt).head
          SessionzierNonPIIKafkaOutputRecord(
            "", topic, nonPIIRecord.getMetadata.getSessionEventTs, nonPIIRecord)
        })
        .sinkTo(getKafkaSessionzierNonPIISink)
        .name("sessionizer_out_nonpii_kafka_sink")
    }
  }

  private def setRestartStrategy = {
    env.setRestartStrategy(RestartStrategies.fixedDelayRestart(
      3, // number of restart attempts
      Time.seconds(10).toMilliseconds // delay
    ))
  }

  def build: JobFlowGraph = {
    env.getConfig.setGlobalJobParameters(new MapBasedJobParameters(config.toMap))
    env.setMaxParallelism(config.get(MAX_PARALLELISM))
    setRestartStrategy
    sinkToKafka(config.get(DV_DEMAND_PARTNER_INTERACTION_TOPIC))
    sinkToKafka(config.get(DV_AUCTION_TOPIC))
    sinkToKafka(config.get(SESSIONIZER_OUT_TOPIC))
    sinkToKafka(config.get(SESSIONIZER_OUT_NON_PII_TOPIC))
    this
  }

  def execute = env.execute(config.get(JOB_NAME))

}

问题分析与优化方案

并行度差异的核心原因

  1. 多Source实例导致资源浪费
    你的代码为每个Sink链路创建了独立的Kafka Source实例,这意味着作业中存在4个完全相同的Source算子同时消费同一个Topic。这种设计不仅会导致重复消费,还会占用额外的CPU、内存和Kafka连接资源,反而挤压下游算子的处理能力,最终导致输出速率跟不上消费速率。

  2. 下游算子并行度未匹配负载
    虽然Source并行度与Partition数量匹配,但下游的KeyedProcessFunction、FlatMap等算子的并行度默认继承全局配置,如果全局并行度未根据数据处理量调整,或者存在数据倾斜,会成为处理瓶颈:

  • KeyedProcessFunction依赖sessionId做KeyBy,如果sessionId分布不均,会出现部分并行实例负载过高,拖慢整体处理速度。
  • FlatMap算子可能将单条输入扩展为多条输出,数据膨胀会进一步加剧下游压力。
  1. Exactly-Once语义的额外开销
    所有Sink都配置了EXACTLY_ONCE交付语义,事务机制会引入提交、协调等额外开销,如果事务参数配置不合理(如超时时间、批量大小),会限制输出速率。

具体优化措施

  1. 复用单个Kafka Source实例
    重构代码,只创建一个Source流,然后分流到各个处理链路,避免重复消费和资源浪费:
// 全局唯一Source实例
private val sourceStream: DataStream[StreamingRecord] = env.fromSource(getKafkaSource,
  WatermarkStrategy
    .forBoundedOutOfOrderness[StreamingRecord](Duration.ofMinutes(1))
    .withTimestampAssigner(new SerializableTimestampAssigner[StreamingRecord] {
      override def extractTimestamp(element: StreamingRecord, recordTimestamp: Long): Long = element.eventTs
    })
    .withIdleness(Duration.ofMinutes(10))
, "Kafka Source")(streamingRecordType)
相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.01 04:40:59