【问题标题】:Processor API: bulk POST request for events stored in KeyValueStore处理器 API:存储在 KeyValueStore 中的事件的批量 POST 请求
【发布时间】:2020-04-01 14:41:29
【问题描述】:

正如此处https://stackoverflow.com/a/60942154/1690657 所建议的那样,我已使用处理器 API 将传入请求存储在 KeyValueStore 中。每 100 个事件我想发送一个 POST 请求。所以我这样做了:

public class BulkProcessor implements Processor<byte[], UserEvent> {

    private KeyValueStore<Integer, ArrayList<UserEvent>> keyValueStore;

    private BulkAPIClient bulkClient;

    private String storeName;

    private ProcessorContext context;

    private int count;

    @Autowired
    public BulkProcessor(String storeName, BulkClient bulkClient) {
        this.storeName = storeName;
        this.bulkClient = bulkClient;
    }

    @Override
    public void init(ProcessorContext context) {
        this.context = context;
        keyValueStore = (KeyValueStore<Integer, ArrayList<UserEvent>>) context.getStateStore(storeName);
        count = 0;
        // to check every 15 minutes if there are any remainders in the store that are not sent yet
        this.context.schedule(Duration.ofMinutes(15), PunctuationType.WALL_CLOCK_TIME, (timestamp) -> {
            if (count > 0) {
                sendEntriesFromStore();
            }
        });
    }

    @Override
    public void process(byte[] key, UserEvent value) {
        int userGroupId = Integer.valueOf(value.getUserGroupId());
        ArrayList<UserEvent> userEventArrayList = keyValueStore.get(userGroupId);
        if (userEventArrayList == null) {
            userEventArrayList = new ArrayList<>();
        }
        userEventArrayList.add(value);
        keyValueStore.put(userGroupId, userEventArrayList);
        if (count == 100) {
            sendEntriesFromStore();
        }
    }

    private void sendEntriesFromStore() {
        KeyValueIterator<Integer, ArrayList<UserEvent>> iterator = keyValueStore.all();
        while (iterator.hasNext()) {
            KeyValue<Integer, ArrayList<UserEvent>> entry = iterator.next();
            keyValueStore.delete(entry.key);
            BulkRequest bulkRequest = new BulkRequest(entry.key, entry.value);
            if (bulkRequest.getLocation() != null) {
                URI url = bulkClient.buildURIPath(bulkRequest);
                bulkClient.postRequestBulkApi(url, bulkRequest);
            }
        }
        iterator.close();
        count = 0;
    }

    @Override
    public void close() {
    }
}

我不确定添加count 是否是线程安全的,或者这是否是实现它的正确方法。目前我也只从一个分区读取。 所以我的问题是:

  • 这是线程安全的吗?
  • 这是在处理器 API 中发送批量 POST 请求的好方法吗?

【问题讨论】:

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


    【解决方案1】:

    这是线程安全的吗?

    是的,它是线程安全的。每个线程都使用ProcessorSupplier 创建自己的Processor 实例/对象。

    这是在处理器 API 中发送批量 POST 请求的好方法吗?

    我觉得总体不错。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2020-01-10
      • 2023-02-09
      • 2020-09-23
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2016-08-04
      • 2018-12-11
      相关资源
      最近更新 更多