You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Apache Beam:并行下载Google Cloud Storage对象时保持对象分组

解决方案:在Beam中按GCS Glob分组并行读取文件内容

针对你的核心需求——为每个输入的GCS Glob模式,将匹配的所有文件内容聚合为单个Iterable<String>元素,同时保证并行读取与分组可靠性,以下是几种符合Beam设计理念的方案,按推荐优先级排序:

方案一:官方FileIO + Combine.perKey(最优选择)

这个方案完全基于Beam官方组件,稳定可靠,无需自定义复杂逻辑,并行度由Beam自动调度,能很好平衡延迟与正确性。

实现步骤

  1. 为每个输入的Glob添加自身作为Key,确保后续分组能关联回原始Glob。
  2. 使用FileIO.matchAll匹配每个Glob对应的所有文件。
  3. 读取文件内容并保留原始Glob作为Key。
  4. 用Combine.perKey按Glob聚合所有文件内容,最终转换为Iterable<String>输出。

代码示例

import org.apache.beam.sdk.Pipeline;
import org.apache.beam.sdk.io.FileIO;
import org.apache.beam.sdk.io.ReadableFile;
import org.apache.beam.sdk.transforms.*;
import org.apache.beam.sdk.values.KV;
import org.apache.commons.io.IOUtils;

import java.io.IOException;
import java.io.InputStream;
import java.nio.charset.StandardCharsets;

public class GlobToFileContents {
    public static void main(String[] args) {
        Pipeline pipeline = Pipeline.create();
        
        // 输入:每个元素是一个GCS Glob模式
        PCollection<String> inputGlobs = pipeline.apply(Create.of(
            "gs://my-bucket/path/*.txt",
            "gs://my-bucket/another-path/**/*.csv"
        ));

        PCollection<Iterable<String>> result = inputGlobs
            // 为每个Glob添加自身作为Key
            .apply(WithKeys.of(SerializableFunctions.identity()))
            // 匹配每个Key对应的Glob文件
            .apply(FileIO.matchAll())
            // 读取文件内容,保留原始Glob作为Key
            .apply(MapElements.via((KV<String, FileIO.ReadableFile> kv) -> {
                try (InputStream in = kv.getValue().open()) {
                    String content = IOUtils.toString(in, StandardCharsets.UTF_8);
                    return KV.of(kv.getKey(), content);
                } catch (IOException e) {
                    throw new RuntimeException("读取文件失败: " + kv.getValue().getMetadata().resourceId(), e);
                }
            }))
            // 按Glob聚合所有文件内容为列表
            .apply(Combine.perKey(ToList.combineFn()))
            // 去掉Key,得到最终的Iterable<String>
            .apply(Values.create());

        pipeline.run();
    }
}

关键优势

  • 正确性保证:FileIO.matchAll会枚举每个Glob对应的所有文件,Combine.perKey会等待同Key的所有文件内容到达后再聚合,完全避免窗口机制的概率性问题。
  • 低延迟:Beam会自动并行处理文件匹配与读取,无需手动管理异步逻辑。
  • 易维护:基于官方组件,无需自定义复杂的DoFn或状态管理。

方案二:AsyncIO异步读取 + 聚合

如果需要更极致的IO并行度,可以使用Beam官方的AsyncIO组件,异步发起GCS文件读取请求,再聚合结果。

实现思路

针对每个输入的Glob,先同步获取所有匹配的文件路径,再异步读取每个文件的内容,最后聚合为Iterable<String>输出。

代码示例

import org.apache.beam.sdk.transforms.AsyncIO;
import org.apache.beam.sdk.io.FileSystems;
import org.apache.beam.sdk.io.fs.MatchResult;
import org.apache.beam.sdk.io.fs.ResourceId;

import java.io.IOException;
import java.io.InputStream;
import java.util.List;
import java.util.concurrent.CompletableFuture;
import java.util.stream.Collectors;

PCollection<Iterable<String>> result = inputGlobs.apply(AsyncIO.<String, Iterable<String>>create()
    .asyncFn(glob -> {
        // 第一步:获取Glob匹配的所有文件路径
        MatchResult matchResult;
        try {
            matchResult = FileSystems.match(glob);
        } catch (IOException e) {
            return CompletableFuture.failedFuture(new RuntimeException("匹配Glob失败: " + glob, e));
        }

        List<String> filePaths = matchResult.metadata().stream()
            .map(MatchResult.Metadata::resourceId)
            .map(ResourceId::toString)
            .collect(Collectors.toList());

        // 第二步:异步读取每个文件内容
        List<CompletableFuture<String>> readFutures = filePaths.stream()
            .map(path -> CompletableFuture.supplyAsync(() -> {
                try (InputStream in = FileSystems.open(ResourceId.fromString(path))) {
                    return IOUtils.toString(in, StandardCharsets.UTF_8);
                } catch (IOException e) {
                    throw new RuntimeException("读取文件失败: " + path, e);
                }
            }))
            .collect(Collectors.toList());

        // 第三步:等待所有异步读取完成,聚合结果
        return CompletableFuture.allOf(readFutures.toArray(new CompletableFuture[0]))
            .thenApply(v -> readFutures.stream()
                .map(CompletableFuture::join)
                .collect(Collectors.toList()));
    })
    .outputType(TypeDescriptor.of(new TypeToken<Iterable<String>>() {})));

关键优势

  • 高并行IO:异步读取能充分利用网络带宽,适合文件数量多、单文件小的场景。
  • 符合Beam设计:比手动使用AsyncHttpClient更规范,Beam会自动管理异步任务的生命周期与重试。

方案三:自定义SplittableDoFn(精细控制场景)

如果需要对文件读取过程做高度自定义的控制(比如特殊的重试逻辑、进度跟踪),可以实现自定义的SplittableDoFn(SDF),将单个Glob的文件读取拆分为多个并行子任务,最后聚合结果。

核心实现要点

  1. 定义Restriction:创建GlobRestriction类,存储Glob模式与待处理的文件路径列表。
  2. 拆分Restriction:在SplitRestriction方法中将文件列表拆分为多个子列表,实现并行处理。
  3. 状态聚合:使用BagState收集所有文件内容,当所有子任务完成后输出聚合结果。

简化代码示例

import org.apache.beam.sdk.state.BagState;
import org.apache.beam.sdk.state.StateSpec;
import org.apache.beam.sdk.state.StateSpecs;
import org.apache.beam.sdk.transforms.splittabledofn.*;

public class GlobToFileContentsSDF extends SplittableDoFn<String, Iterable<String>> {

    @StateId("content")
    private final StateSpec<BagState<String>> contentState = StateSpecs.bag();

    // 定义Restriction:存储Glob与待处理文件路径
    private static class GlobRestriction implements Restriction {
        final String glob;
        final List<String> remainingFiles;

        GlobRestriction(String glob, List<String> remainingFiles) {
            this.glob = glob;
            this.remainingFiles = remainingFiles;
        }
    }

    // 实现RestrictionTracker:管理文件处理进度
    private static class GlobTracker implements RestrictionTracker<GlobRestriction, String> {
        private final GlobRestriction restriction;
        private int nextIndex = 0;

        GlobTracker(GlobRestriction restriction) {
            this.restriction = restriction;
        }

        @Override
        public boolean tryClaim(String filePath) {
            if (nextIndex < restriction.remainingFiles.size()) {
                nextIndex++;
                return true;
            }
            return false;
        }

        @Override
        public GlobRestriction currentRestriction() {
            return new GlobRestriction(restriction.glob,
                restriction.remainingFiles.subList(nextIndex, restriction.remainingFiles.size()));
        }

        @Override
        public SplitResult<GlobRestriction> split(double fraction) {
            int splitPoint = (int) (nextIndex + (restriction.remainingFiles.size() - nextIndex) * fraction);
            if (splitPoint <= nextIndex || splitPoint >= restriction.remainingFiles.size()) {
                return null;
            }
            return SplitResult.of(
                new GlobRestriction(restriction.glob, restriction.remainingFiles.subList(nextIndex, splitPoint)),
                new GlobRestriction(restriction.glob, restriction.remainingFiles.subList(splitPoint, restriction.remainingFiles.size()))
            );
        }

        @Override
        public IsBounded isBounded() { return IsBounded.BOUNDED; }
        @Override
        public void checkDone() {}
    }

    @Override
    public CreateRestriction<String, GlobRestriction> createRestriction() {
        return input -> {
            MatchResult matchResult = FileSystems.match(input);
            List<String> filePaths = matchResult.metadata().stream()
                .map(m -> m.resourceId().toString())
                .collect(Collectors.toList());
            return new GlobRestriction(input, filePaths);
        };
    }

    @Override
    public SplitRestriction<String, GlobRestriction> splitRestriction() {
        return (element, restriction) -> {
            // 按固定大小拆分文件列表,实现并行处理
            int chunkSize = 10;
            List<GlobRestriction> splits = new ArrayList<>();
            for (int i = 0; i < restriction.remainingFiles.size(); i += chunkSize) {
                int end = Math.min(i + chunkSize, restriction.remainingFiles.size());
                splits.add(new GlobRestriction(restriction.glob, restriction.remainingFiles.subList(i, end)));
            }
            return splits;
        };
    }

    @Override
    public RestrictionTracker<GlobRestriction, String> newTracker(GlobRestriction restriction) {
        return new GlobTracker(restriction);
    }

    @ProcessElement
    public void processElement(@Element String glob,
                               RestrictionTracker<GlobRestriction, String> tracker,
                               OutputReceiver<Iterable<String>> receiver,
                               @StateId("content") BagState<String> contentState) throws IOException {
        GlobRestriction restriction = tracker.currentRestriction();
        for (String filePath : restriction.remainingFiles) {
            if (!tracker.tryClaim(filePath)) break;
            // 读取文件内容并存入状态
            try (InputStream in = FileSystems.open(ResourceId.fromString(filePath))) {
                contentState.add(IOUtils.toString(in, StandardCharsets.UTF_8));
            }
        }
        // 所有文件处理完成后,输出聚合结果
        if (tracker.isDone()) {
            receiver.output(contentState.read());
            contentState.clear();
        }
    }
}

关键优势

  • 精细控制:可以自定义文件拆分规则、读取逻辑、进度跟踪等。
  • 原生并行:符合Beam的SDF模型,能自动适配集群的并行能力。

内容的提问来源于stack exchange,提问作者egalpin

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.04.29 12:39:11