Apache Beam:DefaultFilenamePolicy参数类型与文档示例不符问题
解决Apache Beam DynamicAvroDestinations中DefaultFilenamePolicy的String路径问题
我之前也踩过这个坑!确实,Beam 2.4.0的AvroIO官方文档里的这个示例有点误导人——DefaultFilenamePolicy.Params.withBaseFilename()方法并不支持直接传入String类型的路径,它只接受ResourceId对象。
解决办法:把String路径转为ResourceId
你需要借助org.apache.beam.sdk.io.FileSystems类的matchNewResource()方法,将拼接好的字符串路径转换成ResourceId实例,再传给withBaseFilename()。
修改后的代码示例如下:
import org.apache.beam.sdk.io.FileSystems; import org.apache.beam.sdk.io.fs.ResourceId; import java.io.IOException; public FilenamePolicy getFilenamePolicy(Integer userId) throws IOException { String basePath = baseDir + "/user-" + userId + "/events"; // 将String路径转为ResourceId,第二个参数表示是否是目录 ResourceId baseResourceId = FileSystems.matchNewResource(basePath, true); return DefaultFilenamePolicy.fromParams( new Params().withBaseFilename(baseResourceId) ); }
额外注意点
- 别忘了处理
IOException,因为matchNewResource()会抛出这个检查型异常,你可以选择在方法上声明抛出,或者用try-catch包裹处理。 - 第二个参数
isDirectory要根据你的实际场景设置:如果你的baseFilename指向的是目录,就传true;如果是具体文件前缀,传false。
这个问题大概率是文档更新不及时导致的——文档示例可能沿用了更早版本的API写法,但实际实现已经调整为要求ResourceId类型了。
内容的提问来源于stack exchange,提问作者bjorndv
相关产品推荐
相关产品推荐

