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

求助:Spring Batch使用BigQuery替代RDBMS作为ItemReader及分页方案

Spring Batch集成BigQuery读取数据并写入Pub/Sub实现方案

一、核心集成思路

Spring Batch原生没有提供BigQuery专属的ItemReader,可以基于Google BigQuery Java客户端库自定义实现,结合两种分页方案满足不同场景需求,再配合Pub/Sub客户端实现数据写入。

二、分页读取方案

1. 游标分页(通用场景)

利用BigQuery查询返回的pageToken实现分页,适合无有序主键的表,每次读取一页数据,直到无更多结果。

自定义ItemReader实现

@Component
@StepScope
public class BigQueryCursorItemReader implements ItemReader<User> {
    private final BigQuery bigQuery;
    private String pageToken;
    private Iterator<FieldValueList> currentPageIterator;
    // 替换为你的BQ查询语句
    private static final String BQ_QUERY = "SELECT id, name, email FROM `your-project.your-dataset.your-table`";
    // 单页读取量,根据内存情况调整
    private static final int PAGE_SIZE = 1000;

    public BigQueryCursorItemReader(BigQuery bigQuery) {
        this.bigQuery = bigQuery;
    }

    @Override
    public User read() throws Exception {
        // 当当前页数据耗尽时,请求下一页
        if (currentPageIterator == null || !currentPageIterator.hasNext()) {
            QueryJobConfiguration queryConfig = QueryJobConfiguration.newBuilder(BQ_QUERY)
                    .setPageSize(PAGE_SIZE)
                    .setPageToken(pageToken)
                    .build();

            Job queryJob = bigQuery.create(JobInfo.of(queryConfig));
            // 等待查询完成
            queryJob = queryJob.waitFor();

            if (queryJob == null || !queryJob.isDone()) {
                throw new RuntimeException("BigQuery查询作业执行失败或超时");
            }

            QueryResult queryResult = queryJob.getQueryResults();
            // 更新下一页的token
            pageToken = queryResult.getNextPageToken();
            currentPageIterator = queryResult.iterateAll().iterator();

            // 无更多数据时返回null,结束读取
            if (!currentPageIterator.hasNext()) {
                return null;
            }
        }

        // 转换BQ行数据为业务对象
        FieldValueList row = currentPageIterator.next();
        return User.builder()
                .id(row.get("id").getLongValue())
                .name(row.get("name").getStringValue())
                .email(row.get("email").getStringValue())
                .build();
    }
}

2. 范围分页(适合有序主键场景)

如果表存在自增ID、时间戳等有序字段,可按字段范围分段读取,避免游标分页的状态依赖,更适合大规模数据并行处理。

自定义ItemReader实现

@Component
@StepScope
public class BigQueryRangeItemReader implements ItemReader<User> {
    private final BigQuery bigQuery;
    private long currentStartId = 0;
    // 替换为你的BQ表范围查询模板
    private static final String QUERY_TEMPLATE = "SELECT id, name, email FROM `your-project.your-dataset.your-table` WHERE id > ? AND id <= ? ORDER BY id";
    // 每段读取量
    private static final int BATCH_SIZE = 1000;
    private Iterator<FieldValueList> currentBatchIterator;
    private final long maxId;

    public BigQueryRangeItemReader(BigQuery bigQuery) throws Exception {
        this.bigQuery = bigQuery;
        // 预先查询表中最大ID,作为分段终点
        QueryResult maxIdResult = bigQuery.query(QueryJobConfiguration.newBuilder(
                "SELECT MAX(id) as max_id FROM `your-project.your-dataset.your-table`"
        ).build());
        this.maxId = maxIdResult.iterateAll().iterator().next().get("max_id").getLongValue();
    }

    @Override
    public User read() throws Exception {
        if (currentBatchIterator == null || !currentBatchIterator.hasNext()) {
            long endId = Math.min(currentStartId + BATCH_SIZE, maxId);
            // 所有数据读取完成,返回null
            if (currentStartId >= maxId) {
                return null;
            }

            QueryJobConfiguration queryConfig = QueryJobConfiguration.newBuilder(
                    String.format(QUERY_TEMPLATE, currentStartId, endId)
            ).build();

            Job queryJob = bigQuery.create(JobInfo.of(queryConfig));
            queryJob = queryJob.waitFor();

            if (queryJob == null || !queryJob.isDone()) {
                throw new RuntimeException("BigQuery查询作业执行失败或超时");
            }

            QueryResult queryResult = queryJob.getQueryResults();
            currentBatchIterator = queryResult.iterateAll().iterator();
            // 更新下一段的起始ID
            currentStartId = endId;

            if (!currentBatchIterator.hasNext()) {
                return null;
            }
        }

        FieldValueList row = currentBatchIterator.next();
        return User.builder()
                .id(row.get("id").getLongValue())
                .name(row.get("name").getStringValue())
                .email(row.get("email").getStringValue())
                .build();
    }
}

三、Pub/Sub写入实现

利用Spring Cloud GCP的Pub/Sub模板实现ItemWriter,将读取到的数据序列化后发送到指定Topic:

@Component
public class PubSubItemWriter implements ItemWriter<User> {
    private final PubSubTemplate pubSubTemplate;
    // 替换为你的Pub/Sub Topic完整路径
    private static final String TARGET_TOPIC = "projects/your-project/topics/your-topic";
    private final ObjectMapper objectMapper = new ObjectMapper();

    public PubSubItemWriter(PubSubTemplate pubSubTemplate) {
        this.pubSubTemplate = pubSubTemplate;
    }

    @Override
    public void write(List<? extends User> items) throws Exception {
        // 批量发送数据到Pub/Sub
        for (User user : items) {
            String payload = objectMapper.writeValueAsString(user);
            pubSubTemplate.publish(TARGET_TOPIC, payload);
        }
    }
}

四、Batch作业配置

将自定义的Reader和Writer组装成Step和Job:

@Configuration
@EnableBatchProcessing
public class BqToPubSubJobConfig {
    private final JobBuilderFactory jobBuilderFactory;
    private final StepBuilderFactory stepBuilderFactory;
    private final BigQueryCursorItemReader bqItemReader;
    private final PubSubItemWriter pubSubItemWriter;

    public BqToPubSubJobConfig(JobBuilderFactory jobBuilderFactory,
                               StepBuilderFactory stepBuilderFactory,
                               BigQueryCursorItemReader bqItemReader,
                               PubSubItemWriter pubSubItemWriter) {
        this.jobBuilderFactory = jobBuilderFactory;
        this.stepBuilderFactory = stepBuilderFactory;
        this.bqItemReader = bqItemReader;
        this.pubSubItemWriter = pubSubItemWriter;
    }

    @Bean
    public Step bqToPubSubStep() {
        return stepBuilderFactory.get("bqToPubSubStep")
                .<User, User>chunk(1000) // Chunk大小建议与Reader的分页大小一致
                .reader(bqItemReader)
                .writer(pubSubItemWriter)
                .build();
    }

    @Bean
    public Job bqToPubSubJob() {
        return jobBuilderFactory.get("bqToPubSubJob")
                .start(bqToPubSubStep())
                .build();
    }
}

五、关键注意事项

  • 依赖配置:需引入google-cloud-bigquery、spring-cloud-gcp-starter-pubsub、spring-boot-starter-batch等核心依赖
  • 凭证管理:通过环境变量GOOGLE_APPLICATION_CREDENTIALS指定GCP服务账号JSON文件路径,实现BigQuery和Pub/Sub的授权访问
  • 性能优化:根据服务器内存和数据量调整分页/Chunk大小;大规模数据场景下,可结合Spring Batch的分区(Partitioning)机制实现多线程并行读取
  • 容错处理:可给Reader添加重试逻辑处理BigQuery临时错误;若需支持作业重启,游标分页需将pageToken持久化到作业执行上下文

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.12 02:25:00