【问题标题】:Implement micro batching in java在java中实现微批处理
【发布时间】:2018-05-06 03:26:42
【问题描述】:

我正在开发一个基于 kafka 的应用程序,其中 kafka 侦听器将侦听记录;一旦 kafka 接收到记录,我可能需要将记录写入文件。 这里要将记录写入文件,我们要使用带有批处理大小和超时设置的微批处理。 例如,batchsize 为 10,超时设置为 1000 ms,这意味着在写入文件之前等待 10 条记录,等待时间为 1000 毫秒。如果在任何情况下 Kafka 在 1000 毫秒内仅收到 5 条记录,则在该批次中只写入 5 条记录。

我在 Java 中做到这一点的效率如何。

【问题讨论】:

    标签: java apache-kafka kafka-consumer-api


    【解决方案1】:

    在这种情况下,一种常见的方法是将所有记录放入队列中。当队列达到 10 或 1000 毫秒后,有一个线程将记录这些记录,具体取决于先出现的情况。

    消费者代码:

     CountDownLatch countDownLatch = new CountDownLatch(10);
     countDownLatch.await(1000, TimeUnit.MILLISECONDS);
     int queueSize = queue.size();
     for(int i = 0; i < queueSize; ++i) {
         ... do your work here or put in a batch a do it right after loop
     }
    

    生产者代码:

     Record record = ...receive new record...
     queue.put(record);
     consumer.getCountDownLatch().countDown();
    

    作为队列,我推荐使用无绑定队列,例如LinkedTransferQueue,因为您不想在达到 10 个任务时停止生产者,您仍然需要使用来自 kafka 的结果。

    另外一个选项是reactive streams

    【讨论】:

      【解决方案2】:

      听起来您应该使用Kafka Connect API。这是part of Apache Kafka,旨在支持您描述的那种流程。

      有一个developer guide here

      【讨论】:

        猜你喜欢
        • 2012-09-17
        • 1970-01-01
        • 1970-01-01
        • 2016-06-25
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2014-02-04
        • 1970-01-01
        相关资源
        最近更新 更多