Apache Beam基于日期过滤CSV数据及读取表头、访问列的实现咨询
问题:Apache Beam 读取CSV并按日期过滤的实现疑问
我正在尝试从CSV文件读取记录并基于日期值过滤数据,已按照下述方式完成实现,目前可得到正确结果,但不确定该实现方式是否合规。同时我存在两个核心问题:
- 如何访问CSV的单独列,以便基于对应列的数值设置过滤条件?
- 读取CSV数据时是否可以指定表头?
实现步骤
- 创建Pipeline
- 从文件读取数据
- 执行所需过滤操作
- 创建MapElement对象将OrderRequest转换为String
- 将OrderRequest实体映射为String
- 将输出结果写入文件
实现代码
// Creating pipeline Pipeline pipeline = Pipeline.create(); // For transformations Reading from a file PCollection<String> orderRequest = pipeline .apply(TextIO.read().from("src/main/resources/ST/STCheck/OrderRequest.csv")); PCollection<OrderRequest> pCollectionTransformation = orderRequest .apply(ParDo.of(new DoFn<String, OrderRequest>() { private static final long serialVersionUID = 1L; @ProcessElement public void processElement(ProcessContext c) { String rowString = c.element(); if (!rowString.contains("order_id")) { String[] strArr = rowString.split(","); OrderRequest orderRequest = new OrderRequest(); orderRequest.setOrder_id(strArr[0]); // Condition to check if the String source1 = strArr[1]; DateTimeFormatter fmt1 = DateTimeFormat.forPattern("mm/dd/yyyy"); DateTime d1 = fmt1.parseDateTime(source1); System.out.println(d1); String source2 = "4/24/2017"; DateTimeFormatter fmt2 = DateTimeFormat.forPattern("mm/dd/yyyy"); DateTime d2 = fmt2.parseDateTime(source2); System.out.println(d2); orderRequest.setOrder_date(strArr[1]); System.out.println(strArr[1]); orderRequest.setAmount(Double.valueOf(strArr[2])); orderRequest.setCounter_id(strArr[3]); if (DateTimeComparator.getInstance().compare(d1, d2) > -1) { c.output(orderRequest); } } } })); // Create a MapElement Object and convert the OrderRequest to String MapElements<OrderRequest, String> mapElements = MapElements.into(TypeDescriptors.strings()) .via((OrderRequest orderRequestType) -> orderRequestType.getOrder_id() + " " + orderRequestType.getOrder_date() + " " + orderRequestType.getAmount() + " " + orderRequestType.getCounter_id()); // Mapping the OrderRequest Entity to String PCollection<String> pStringList = pCollectionTransformation.apply(mapElements); // Now Writing the elements to a file pStringList.apply(TextIO.write().to("src/main/resources/ST/STCheck/OrderRequestOut.csv").withNumShards(1) .withSuffix(".csv")); // To run pipeline pipeline.run(); System.out.println("We are done!!");
POJO类定义
public class OrderRequest implements Serializable{ String order_id; String order_date; double amount; String counter_id; // 省略getter、setter方法 }
解答
问题1:如何访问CSV的单独列设置过滤条件
你当前通过split拆分字符串后按索引取字段的方式是可运行的,更规范的方式是使用Beam内置的CsvIO,不需要手动处理字符串拆分、特殊字符转义等逻辑,读取后可以直接通过列名访问对应字段,设置过滤条件更清晰也更稳定。
问题2:读取CSV时是否可以指定表头
可以指定表头,有两种常用方案:
- 如果CSV文件第一行本身就是表头,调用
CsvIO.read().withHeader()即可自动识别表头,不需要你手动写逻辑跳过首行 - 如果CSV文件没有自带表头,你可以通过
CsvIO.read().withHeader("order_id", "order_date", "amount", "counter_id")的方式手动自定义表头,读取后直接用列名获取对应字段即可。
现有实现优化建议
你当前的实现可以正常得到结果,但存在几个可优化的点:
- 你通过
rowString.contains("order_id")跳过表头的逻辑不够严谨,若某行数据本身包含order_id字符串会被误过滤,用内置IO的表头处理逻辑更稳妥 - 日期格式化对象
DateTimeFormatter可以放到DoFn的@Setup方法中初始化,避免每条数据都重复创建对象,提升运行性能 - 直接用
split(",")拆分CSV如果遇到字段内容本身带逗号的场景会拆分错误,内置的CSV解析会自动处理转义场景,容错性更高
内容的提问来源于stack exchange,提问作者User27854
相关产品推荐
相关产品推荐

