You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何将按配置batchSize分批发Kafka的Java7代码改写为Java8风格

Java8 批量遍历列表发送Kafka实现方案

前置依赖类定义(无Lombok版本)

MetaData 类

public class MetaData {
    private int batchNo;
    private int totalSize;
    private int batchSize;

    public MetaData(int batchNo, int totalSize, int batchSize) {
        this.batchNo = batchNo;
        this.totalSize = totalSize;
        this.batchSize = batchSize;
    }

    public int getBatchNo() {
        return batchNo;
    }

    public void setBatchNo(int batchNo) {
        this.batchNo = batchNo;
    }

    public int getTotalSize() {
        return totalSize;
    }

    public void setTotalSize(int totalSize) {
        this.totalSize = totalSize;
    }

    public int getBatchSize() {
        return batchSize;
    }

    public void setBatchSize(int batchSize) {
        this.batchSize = batchSize;
    }
}

DataResponseBean 类

public class DataResponseBean {
    private MetaData metaData;
    private List<String> dataList;

    public DataResponseBean(MetaData metaData, List<String> dataList) {
        this.metaData = metaData;
        this.dataList = dataList;
    }

    public MetaData getMetaData() {
        return metaData;
    }

    public void setMetaData(MetaData metaData) {
        this.metaData = metaData;
    }

    public List<String> getDataList() {
        return dataList;
    }

    public void setDataList(List<String> dataList) {
        this.dataList = dataList;
    }
}

核心批量发送逻辑

public void batchSendToKafka(List<String> responseList, int batchSize) {
    int totalSize = responseList.size();
    // 计算总批次数,向上取整
    int batchCount = (totalSize + batchSize - 1) / batchSize;

    IntStream.range(0, batchCount)
            .forEach(batchIndex -> {
                // 计算当前批次的起止索引
                int from = batchIndex * batchSize;
                int to = Math.min(from + batchSize, totalSize);
                // 截取批次数据,subList为视图无额外拷贝开销
                List<String> currentBatch = responseList.subList(from, to);
                // 构造请求对象
                MetaData meta = new MetaData(batchIndex + 1, totalSize, currentBatch.size());
                DataResponseBean sendBean = new DataResponseBean(meta, currentBatch);
                // 发送kafka,抛出异常时直接终止后续处理
                publishToKafka(sendBean);
            });
}

// 原有Kafka发送方法,异常直接向外抛出即可
private void publishToKafka(DataResponseBean bean) {
    // 原有发送逻辑
}

注意事项

  • 异常处理:publishToKafka 抛出任意异常都会直接终止流遍历,不会继续处理剩余批次,符合中断要求
  • 边界兼容:最后一批元素不足batchSize时会自动取剩余所有元素发送,无需额外判断
  • 数据安全:如果需要修改当前批次数据且不影响原列表,可将subList逻辑替换为new ArrayList<>(responseList.subList(from, to))生成独立副本
  • Lombok适配:后续引入Lombok后,可直接给两个Bean添加@Data、@AllArgsConstructor注解,删除手动编写的构造器、getter、setter代码

内容的提问来源于stack exchange,提问作者Debo

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.10.07 10:54:03