【发布时间】: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