【发布时间】:2020-08-10 19:22:10
【问题描述】:
我有一个包含两列“keyCol”列和“valCol”列的 spark 数据框。数据框非常庞大,将近 1 亿行。 我想以小批量的方式将数据帧写入/生成到 kafka 主题,即每分钟 10000 条记录。 此 spark 作业将每天运行一次,从而创建此数据帧
如何在下面的代码中实现每分钟 10000 条记录的小批量写入,或者请建议是否有更好/有效的方法来实现。
spark_df.foreachPartition(partitions ->{
Producer<String, String> producer= new KafkaProducer<String, String>(allKafkaParamsMapObj);
while (partitions) {
Row row = partitions.next();
producer.send(new ProducerRecord<String, String>("topicName", row.getAs("keyCol"), row.getAs("valCol")), new Callback() {
@Override
public void onCompletion(RecordMetadata recordMetadata, Exception e) {
//Callback code goes here
}
});
}
return;
});
【问题讨论】:
标签: dataframe apache-spark apache-kafka spark-streaming kafka-producer-api