【问题标题】:How to control number of records while writing Spark Dataframe to Kafka Producer using Spark Java如何在使用 Spark Java 将 Spark Dataframe 写入 Kafka Producer 时控制记录数
【发布时间】: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


    【解决方案1】:

    你可以使用下面的grouped(10000)函数并执行一分钟的睡眠线程

    config.foreachPartition(f => {
          f.grouped(10000).foreach( (roqSeq : Seq[Row]) => { // Run 10000 in batch
    
            roqSeq.foreach( row => {
              producer.send(new Nothing("topicName", row.getAs("keyCol"), row.getAs("valCol")), new Nothing() {
                def onCompletion(recordMetadata: Nothing, e: Exception): Unit = {
                  //Callback code goes here
                }
              })
            })
              Thread.sleep(60000) // Sleep for 1 minute
            }
          )
        })
    

    【讨论】:

    • 上述代码无法在 java 中实现,因为 foreachPartition 'f'(类型为 Iterator)在 spark java 中没有 'grouped' 方法。你能建议同样的java代码吗?谢谢
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-06-12
    • 1970-01-01
    • 1970-01-01
    • 2023-02-01
    • 2021-03-06
    • 1970-01-01
    相关资源
    最近更新 更多