如何在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
相关产品推荐
相关产品推荐

