如何基于List<String>配置Text.IO.ReadAll读取GCS多文件夹JSON到Dataflow
如何用Dataflow的Text.IO.ReadAll读取List中的GCS文件夹路径?
嗨,这个问题其实很容易解决——核心就是把你本地的List<String>转换成Dataflow管道能识别的PCollection<String>,然后就能直接用TextIO.readAll()批量处理所有文件夹了。我给你一步步拆解:
第一步:准备好你的GCS文件夹路径列表
先把30个GCS文件夹的路径整理成List<String>,注意GCS路径格式要正确,比如gs://your-bucket-name/target-folder/(末尾加斜杠更稳妥,确保指向文件夹而非同名文件)。第二步:将List转换成PCollection
用Dataflow的Create转换,把本地列表包装成管道内的PCollection<String>,这是连接本地数据和Dataflow管道的关键一步。示例代码如下:import org.apache.beam.sdk.values.PCollection; import org.apache.beam.sdk.transforms.Create; import java.util.Arrays; import java.util.List; // 你的30个GCS文件夹路径列表 List<String> gcsFolderPaths = Arrays.asList( "gs://your-bucket/folder-01/", "gs://your-bucket/folder-02/", // ... 剩下的28个文件夹路径 ); // 把List转换成PCollection PCollection<String> folderPathCollection = pipeline.apply( "Load Folder Paths into Pipeline", Create.of(gcsFolderPaths) );第三步:用TextIO.readAll()批量读取JSON文件
现在你就可以用TextIO.readAll()来读取所有文件夹下的JSON文件了。为了精准只读取JSON格式的文件,建议添加后缀过滤或者文件名过滤:import org.apache.beam.sdk.io.TextIO; PCollection<String> jsonFileContents = folderPathCollection.apply( "Read All JSON Files from Folders", TextIO.readAll() .withSuffix(".json") // 只读取后缀为.json的文件 .fromGlob("*") // 匹配每个文件夹下的所有符合条件的文件 );如果需要更复杂的文件名过滤规则(比如包含特定关键词),可以用
withFilenameFilter:.withFilenameFilter((SerializableFunction<String, Boolean>) filename -> filename.endsWith(".json") && filename.contains("target-keyword") )
这样配置后,Dataflow就会自动遍历你列表里的所有GCS文件夹,把里面的JSON文件以字符串形式读取到管道中,完全不用为每个文件夹单独添加Read步骤~
内容的提问来源于stack exchange,提问作者PUG
相关产品推荐
相关产品推荐

