Apache Beam下PCollection等值连接代码报编译错误问题咨询
问题原因与解决方案
官方文档的Join代码示例本身是正确的,你遇到的编译错误是使用方式不符合API要求导致的,核心问题出在两点:
- 参与Join的PCollection未配置Schema:Beam的
org.apache.beam.sdk.schemas.transforms.Join系列API属于Schema感知的转换,要求输入的PCollection必须附带Schema信息,要么是原生Row类型,要么是自定义类通过@DefaultSchema注解注册了Schema,否则Beam无法解析你指定的连接字段userId、productId对应的元数据,导致泛型推导失败。 - Join调用未显式指定泛型:未指定泛型的情况下Java编译器无法正确推导
Join.innerJoin()返回的PTransform的输入类型,进而不符合apply方法的入参类型要求。
解决步骤
步骤1:为自定义类配置Schema
给你的Transaction和Review类添加@DefaultSchema注解,示例如下:
import org.apache.beam.sdk.schemas.DefaultSchema; import org.apache.beam.sdk.schemas.JavaFieldSchema; @DefaultSchema(JavaFieldSchema.class) public class Transaction { // 字段名要和连接时指定的字段名完全匹配 public String userId; public String productId; // 其他业务字段 // 必须保留无参构造方法 public Transaction() {} } // Review类按同样规则添加注解
如果不允许修改实体类,也可以调用PCollection.setSchema()方法手动给输入PCollection绑定Schema。
步骤2:调整Join调用写法,显式指定泛型
修改你的连接代码,显式声明Join的左右输入类型,保证泛型推导正确:
PCollection<Transaction> transactions = readTransactions(); PCollection<Review> reviews = readReviews(); // 显式指定<Transaction, Review>泛型,匹配输入PCollection的类型 PCollection<Row> joined = transactions.apply( Join.<Transaction, Review>innerJoin(reviews).using("userId", "productId") );
备选方案:使用旧版Join API
如果你使用的是org.apache.beam.sdk.extensions.joinlibrary.Join(非Schema感知的旧版API),需要先将两个PCollection转换为以连接键为key的KV类型再做连接:
// 拼接多字段为唯一连接键 PCollection<KV<String, Transaction>> txnWithKey = transactions.apply( MapElements.into( TypeDescriptors.kvs(TypeDescriptors.strings(), TypeDescriptor.of(Transaction.class)) ).via(txn -> KV.of(txn.getUserId() + "#" + txn.getProductId(), txn)) ); PCollection<KV<String, Review>> reviewWithKey = reviews.apply( MapElements.into( TypeDescriptors.kvs(TypeDescriptors.strings(), TypeDescriptor.of(Review.class)) ).via(review -> KV.of(review.getUserId() + "#" + review.getProductId(), review)) ); // 调用旧版Inner Join PCollection<KV<String, KV<Transaction, Review>>> joined = txnWithKey.apply(Join.innerJoin(reviewWithKey));
内容的提问来源于stack exchange,提问作者Thomas
相关产品推荐
相关产品推荐

