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

如何在Apache Beam Java Dataflow中替换BigQuery指定月份数据

实现仅替换BigQuery表中指定月份数据的方案

针对你的需求,有三种可行的实现方式,你可以根据表结构和业务场景选择:


方案一:先删除旧数据,再追加新数据

这种方式逻辑简单,先通过DML语句删除目标月份的旧数据,再将新数据追加到表中。

步骤1:添加删除旧数据的流水线步骤

在现有流水线的开头,添加执行BigQuery DELETE语句的步骤,根据你的日期字段类型调整筛选条件:

// 假设表中日期字段为DATE类型的`stat_date`,筛选2023年7月的数据
String deleteOldDataQuery = "DELETE FROM `" + writeCdnMediaRequestTable + "` " +
                            "WHERE EXTRACT(YEAR FROM stat_date) = 2023 AND EXTRACT(MONTH FROM stat_date) = 7";

// 如果日期字段是STRING类型(如'2023-07-01'),用以下语句:
// String deleteOldDataQuery = "DELETE FROM `" + writeCdnMediaRequestTable + "` " +
//                             "WHERE stat_date LIKE '2023-07-%'";

pipeline.apply("Delete 2023-07 Old Data", BigQueryIO.executeQuery(deleteOldDataQuery)
        .usingStandardSql()
        .withCreateDisposition(BigQueryIO.Write.CreateDisposition.CREATE_NEVER)
        .withWriteDisposition(BigQueryIO.Write.WriteDisposition.WRITE_APPEND));

步骤2:修改写入配置为追加模式

将原有的WRITE_TRUNCATE改为WRITE_APPEND,确保新数据追加到表中而非清空全表:

.apply("Write Data To BigQuery", BigQueryIO.writeTableRows()
    .to(writeCdnMediaRequestTable)
    .withSchema(cdnDailyRequestSchema)
    .withCreateDisposition(BigQueryIO.Write.CreateDisposition.CREATE_IF_NEEDED)
    .withWriteDisposition(BigQueryIO.Write.WriteDisposition.WRITE_APPEND));

方案二:使用MERGE语句原子替换数据

这种方式更安全,通过临时表+MERGE操作实现原子性替换,避免删除后写入失败导致的数据缺失。

步骤1:将新数据写入临时表

先把生成的新数据写入一个临时表(可通过时间戳命名避免冲突):

// 生成临时表名
String tempTableName = writeCdnMediaRequestTable + "_temp_" + System.currentTimeMillis();

// 生成待写入的数据(复用你原有的转换逻辑)
PCollection<TableRow> cdnDailyRows = pipeline
        .apply("Read from cdn_requests BigQuery", BigQueryIO
                .read(new CdnMediaRequestLogEntity.FromSchemaAndRecord())
                .fromQuery(cdnRequestsQueryString)
                .usingStandardSql())
        .apply("Validate and Filter Cdn Media Request Log Objects", Filter.by(new CdnMediaRequestValidator()))
        .apply("Convert Cdn Logs To Key Value Pairs", ParDo.of(new CdnMediaRequestResponseSizeKeyValuePairConverter()))
        .apply("Sum the Response Sizes By Key", Sum.longsPerKey())
        .apply("Convert To New Daily Requests Objects", ParDo.of(new CdnDailyRequestConverter(projectId, kind)))
        .apply("Convert Cdn Media Request Entities to Big Query Objects", ParDo.of(new BigQueryCdnDailyRequestRowConverter()));

// 写入临时表
cdnDailyRows.apply("Write to Temp Table", BigQueryIO.writeTableRows()
        .to(tempTableName)
        .withSchema(cdnDailyRequestSchema)
        .withCreateDisposition(BigQueryIO.Write.CreateDisposition.CREATE_IF_NEEDED)
        .withWriteDisposition(BigQueryIO.Write.WriteDisposition.WRITE_TRUNCATE));

步骤2:执行MERGE操作替换数据

通过MERGE语句删除目标表中2023-07的旧数据,同时插入临时表的新数据:

// 构建MERGE语句
String mergeQuery = "MERGE INTO `" + writeCdnMediaRequestTable + "` T " +
                    "USING `" + tempTableName + "` S " +
                    "ON FALSE " + // 不匹配任何数据,用于执行DELETE和INSERT
                    "WHEN NOT MATCHED BY SOURCE AND EXTRACT(YEAR FROM T.stat_date) = 2023 AND EXTRACT(MONTH FROM T.stat_date) = 7 THEN DELETE " +
                    "WHEN NOT MATCHED BY TARGET THEN INSERT ROW";

// 执行MERGE
pipeline.apply("Merge Temp Data into Target Table", BigQueryIO.executeQuery(mergeQuery)
        .usingStandardSql()
        .withCreateDisposition(BigQueryIO.Write.CreateDisposition.CREATE_NEVER));

// 可选:删除临时表
String dropTempTableQuery = "DROP TABLE IF EXISTS `" + tempTableName + "`";
pipeline.apply("Drop Temp Table", BigQueryIO.executeQuery(dropTempTableQuery)
        .usingStandardSql());

如果你的表有唯一业务键(比如stat_date+cdn_id),也可以用匹配键的方式更新数据,替换上述MERGE语句:

String mergeQuery = "MERGE INTO `" + writeCdnMediaRequestTable + "` T " +
                    "USING `" + tempTableName + "` S " +
                    "ON T.stat_date = S.stat_date AND T.cdn_id = S.cdn_id " + // 匹配唯一键
                    "WHEN MATCHED THEN UPDATE SET T.request_count = S.request_count, T.response_size = S.response_size " + // 更新字段
                    "WHEN NOT MATCHED AND EXTRACT(YEAR FROM S.stat_date) = 2023 AND EXTRACT(MONTH FROM S.stat_date) = 7 THEN INSERT ROW";

方案三:利用分区表特性(最优方案)

如果你的表是按日期分区的(推荐改造为分区表),可以直接指定分区写入,此时WRITE_TRUNCATE只会清空目标分区而非全表,操作更高效。

步骤1:确保表是日期分区表

如果还未分区,可通过BigQuery控制台或DDL语句将表改为按stat_date(DATE类型)分区。

步骤2:修改写入配置指定分区

在写入时指定目标分区,保留WRITE_TRUNCATE仅清空该分区:

.apply("Write Data To BigQuery", BigQueryIO.writeTableRows()
    .to(new TableDestination(writeCdnMediaRequestTable, "Write 2023-07 partition")
            .setPartitioning(Partitioning.newBuilder()
                    .setField("stat_date")
                    .setType(Partitioning.Type.DATE)
                    .build()))
    .withSchema(cdnDailyRequestSchema)
    .withCreateDisposition(BigQueryIO.Write.CreateDisposition.CREATE_IF_NEEDED)
    .withWriteDisposition(BigQueryIO.Write.WriteDisposition.WRITE_TRUNCATE));

或者直接指定分区标识符(适用于已存在的分区表):

.apply("Write Data To BigQuery", BigQueryIO.writeTableRows()
    .to(TableReference.of(writeCdnMediaRequestTable).setPartition("202307")) // 分区格式为YYYYMM
    .withSchema(cdnDailyRequestSchema)
    .withCreateDisposition(BigQueryIO.Write.CreateDisposition.CREATE_IF_NEEDED)
    .withWriteDisposition(BigQueryIO.Write.WriteDisposition.WRITE_TRUNCATE));

内容的提问来源于stack exchange,提问作者Rey Marvin De Jesus

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 21:35:57