【问题标题】:How to send data as key/value to Kafka using Apache Flink with Specific Partitioner如何使用带有特定分区器的 Apache Flink 将数据作为键/值发送到 Kafka
【发布时间】:2021-04-11 03:47:07
【问题描述】:

我在 Flink 中有一个有效载荷,如下所示;

{
    "memberId": 4
    "total": 5
}

我想使用指定的分区器将数据作为键值格式发送到 kafka。对于分区器,我将使用 Modulo 分区器。

模分区器示例;

partitionId = value % numPartitions

假设numPartitions参数为3。如果我们可以使用上面定义的payload的memberId,partitionId应该是4 % 3 = 1

根据上面的partitioner,我想将具有相同partitionId的数据发送到同一个kafka topic。另一个例子;

如果(假设 numPartitions = 3);

memberId: 3 => (3 % 3) => partitionId = 0 => kafka partition 1
memberId: 8 => (8 % 3) => partitionId = 2 => kafka partition 2
memberId: 2 => (2 % 3) => partitionId = 2 => kafka partition 2
memberId: 6 => (6 % 3) => partitionId = 0 => kafka partition 1
memberId: 7 => (7 % 3) => partitionId = 1 => kafka partition 2

如果我没记错的话,如果我们不能指定任何 key 和 partition 函数,flink kafka producer 使用 FlinkFixedPartitioner。如果我们将分区函数设置为null,flink kafka producer 将使用循环分发。但我不知道如何将数据作为键/值格式发送到 kafka,如何通过模数对其进行分区。我怎样才能做到这一点?

【问题讨论】:

    标签: apache-kafka apache-flink flink-streaming


    【解决方案1】:

    如果您使用KafkaSerializationSchema,那么您可以创建 Kafka ProducerRecords,并设置 Kafka 键(和值)。也可以在ProducerRecord中设置分区。

    【讨论】:

    • 那么,模分区器呢?
    • 哦,谢谢! ProducerRecord可以带kafka分区!
    猜你喜欢
    • 2019-02-04
    • 1970-01-01
    • 2018-10-23
    • 1970-01-01
    • 2018-04-18
    • 2018-12-03
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多