【发布时间】:2019-03-07 18:15:33
【问题描述】:
我正在使用 Spark 流式传输,并且正在将数据发送到 Kafka。我正在向 Kafka 发送地图。假设我有一个 20 的 Map(在 Streaming Batch 持续时间内可能会增长到 1000)元素,如下所示:
HashMap<Integer,String> input = new HashMap<Integer,String>();
input.put(11,"One");
input.put(312,"two");
input.put(33,"One");
input.put(24,"One");
input.put(35,"One");
input.put(612,"One");
input.put(7,"One");
input.put(128,"One");
input.put(9,"One");
input.put(10,"One");
input.put(11,"One1");
input.put(12,"two1");
input.put(13,"One1");
input.put(14,"One1");
input.put(15,"One1");
input.put(136,"One1");
input.put(137,"One1");
input.put(158,"One1");
input.put(159,"One1");
input.put(120,"One1");
Set<Integer> inputKeys = input.keySet();
Iterator<Integer> inputKeysIterator = inputKeys.iterator();
while (inputKeysIterator.hasNext()) {
Integer key = inputKeysIterator.next();
ProducerRecord<Integer, String> record = new ProducerRecord<Integer, String>(topic,
key%10, input.get(key));
KafkaProducer.send(record);
}
我的 Kafka 主题有 10 个分区。在这里,我调用了 kafkaProducer.send() 20 次,因此调用了 20 次 Kafka。如何批量发送整个数据,即在一个 Kafka 调用中,但我想再次确保每条记录都转到由公式 key%10 驱动的特定分区,如
ProducerRecord 记录 = 新 ProducerRecord(topic, key%10, input.get(key));
我看到的选项:linger.ms=1 可以确保这一点,但延迟为 1 毫秒。 如何避免这种延迟并避免 20 个网络(Kafka)调用或最小化 Kafka 调用?
【问题讨论】:
标签: apache-spark apache-kafka spark-streaming kafka-producer-api