【发布时间】:2017-02-07 01:36:36
【问题描述】:
谁能提供我从 Spark Streaming 向 Kafka 推送记录的示例代码?
【问题讨论】:
标签: apache-spark apache-kafka spark-streaming
谁能提供我从 Spark Streaming 向 Kafka 推送记录的示例代码?
【问题讨论】:
标签: apache-spark apache-kafka spark-streaming
我已经使用 Java 完成了它。您可以在JavaDStream<String> 上使用此函数作为.foreachRDD() 的参数。这不是最好的方法,因为它会为每个 RDD 创建一个 KafkaProducer,您可以使用 KafkaProducers 的“池”来执行此操作,例如 socket example in Spark documentation。
这是我的代码:
public static class KafkaPublisher implements VoidFunction<JavaRDD<String>> {
private static final long serialVersionUID = 1L;
public void call(JavaRDD<String> rdd) throws Exception {
Properties props = new Properties();
props.put("bootstrap.servers", "loca192.168.0.155lhost:9092");
props.put("acks", "1");
props.put("retries", 0);
props.put("batch.size", 16384);
props.put("linger.ms", 1000);
props.put("buffer.memory", 33554432);
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
rdd.foreachPartition(new VoidFunction<Iterator<String>>() {
private static final long serialVersionUID = 1L;
public void call(Iterator<String> partitionOfRecords) throws Exception {
Producer<String, String> producer = new KafkaProducer<>(props);
while(partitionOfRecords.hasNext()) {
producer.send(new ProducerRecord<String, String>("topic", partitionOfRecords.next()));
}
producer.close();
}
});
}
}
【讨论】:
使用 Spark Streaming,您可以使用来自 Kafka 主题的数据。
如果您想将记录发布到 Kafka 主题,可以使用 Kafka Producer [https://cwiki.apache.org/confluence/display/KAFKA/0.8.0+Producer+Example]
或者您可以使用 Kafka Connect 使用多个源连接器将数据发布到 Kafka 主题。[http://www.confluent.io/product/connectors/]
有关 Spark 流和 Kafka 集成的更多信息,请参阅下面的链接。
http://spark.apache.org/docs/latest/streaming-kafka-integration.html
【讨论】: