Flink Kafka Source算子并行度调整效果咨询
Flink作业并行度与输出速率不匹配问题分析
作业DAG结构
- 包含4个输出Sink
- 通过Kafka Source读取单个Topic,Topic的每个Partition已按Key分区
- 执行链路:
- Source -> KeyedProcessFunction -> Map -> Sink-1
- Source -> KeyedProcessFunction -> Map -> FlatMap -> Sink-2
- Source -> KeyedProcessFunction -> Map -> FlatMap -> Sink-3
- 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)) }
问题分析与优化方案
并行度差异的核心原因
多Source实例导致资源浪费
你的代码为每个Sink链路创建了独立的Kafka Source实例,这意味着作业中存在4个完全相同的Source算子同时消费同一个Topic。这种设计不仅会导致重复消费,还会占用额外的CPU、内存和Kafka连接资源,反而挤压下游算子的处理能力,最终导致输出速率跟不上消费速率。下游算子并行度未匹配负载
虽然Source并行度与Partition数量匹配,但下游的KeyedProcessFunction、FlatMap等算子的并行度默认继承全局配置,如果全局并行度未根据数据处理量调整,或者存在数据倾斜,会成为处理瓶颈:
- KeyedProcessFunction依赖
sessionId做KeyBy,如果sessionId分布不均,会出现部分并行实例负载过高,拖慢整体处理速度。 - FlatMap算子可能将单条输入扩展为多条输出,数据膨胀会进一步加剧下游压力。
- Exactly-Once语义的额外开销
所有Sink都配置了EXACTLY_ONCE交付语义,事务机制会引入提交、协调等额外开销,如果事务参数配置不合理(如超时时间、批量大小),会限制输出速率。
具体优化措施
- 复用单个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)
相关产品推荐
相关产品推荐

