Apache Beam基于上一行值更新当前行及填充缺失值的实现方法
问题解决步骤
报错根因说明
- 报错1:
Iterables.toArray仅支持接收普通JavaIterable实例作为参数,你传入的是Beam分布式数据集类型PCollection,二者属于完全不同的抽象,自然不匹配。 - 报错2:
PCollection<Row[]>是分布式的Row数组集合,不是本地内存中的Java数组对象,不能直接赋值给本地Row[]类型的变量,所有数据处理逻辑必须封装在DoFn内部执行。
核心实现调整
所有空值填充逻辑都在GroupByDate类中实现,无需额外开发KV转数组的DoFn,完整实现代码如下:
import java.time.LocalDate; import java.time.format.DateTimeFormatter; import java.util.Comparator; import java.util.List; import java.util.stream.Collectors; import java.util.stream.StreamSupport; import org.apache.beam.sdk.transforms.DoFn; import org.apache.beam.sdk.values.KV; import org.apache.beam.sdk.values.Row; class GroupByDate extends DoFn<KV<String,Iterable<Row>>, Row> { private static final long serialVersionUID = -1345126662309830332L; private final org.apache.beam.sdk.schemas.Schema beamSchema; // 构造函数传入Beam Schema用于构建新的Row对象 public GroupByDate(org.apache.beam.sdk.schemas.Schema beamSchema) { this.beamSchema = beamSchema; } @ProcessElement public void processElement(ProcessContext context) { Iterable<Row> rows = context.element().getValue(); // 1. 同customerId下的所有记录按date升序排序 List<Row> sortedRows = StreamSupport.stream(rows.spliterator(), false) .sorted(Comparator.comparing(row -> LocalDate.parse(row.getString("date"), DateTimeFormatter.ofPattern("yyyy/M/d")))) .collect(Collectors.toList()); Double lastValidAmount = null; // 2. 遍历排序后记录按规则填充空值 for (Row originalRow : sortedRows) { Double currentAmount = originalRow.getDouble("amount"); Double filledAmount; if (currentAmount == null) { // 分组内首行空值填充0,其余空值继承上一行有效值 filledAmount = lastValidAmount == null ? 0.0 : lastValidAmount; } else { filledAmount = currentAmount; } // 3. 构建填充后的新Row输出 Row outputRow = Row.withSchema(beamSchema) .addValue(originalRow.getString("customerId")) .addValue(originalRow.getString("date")) .addValue(filledAmount) .build(); context.output(outputRow); // 更新上一行有效值缓存 lastValidAmount = filledAmount; } } }
主代码调整
- 修正ToKV的分组字段名,你的Schema中客户ID字段为
customerId,不是ID:
ToKV toKV = new ToKV(); toKV.setColumnName1("customerId"); PCollection<KV<String, Row>> kvRows = rows.apply(ParDo.of(toKV)).setCoder(KvCoder.of(StringUtf8Coder.of(), rows.getCoder())); PCollection<KV<String,Iterable<Row>>> groupedKVRows = kvRows.apply(GroupByKey.create());
- 调用GroupByDate时传入构造好的Beam Schema:
org.apache.beam.sdk.schemas.Schema beamSchema = AvroUtils.toBeamSchema(schema); PCollection<Row> outputRow = groupedKVRows .apply(ParDo.of(new GroupByDate(beamSchema))) .setCoder(RowCoder.of(beamSchema));
逻辑验证
- 分组排序后首行amount为空时,
lastValidAmount为null,自动填充0,符合BS:89482的填充要求 - 后续行amount为空时,取上一行缓存的有效值填充,符合BS:89483继承上一行300的要求
内容的提问来源于stack exchange,提问作者User27854
相关产品推荐
相关产品推荐

