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

如何使用Apache Beam Java SDK读取MinIO存储桶中的文件?

使用Apache Beam Java SDK读取MinIO中的文件

嘿,我来帮你搞定这个问题!MinIO完全兼容Amazon S3的API,所以我们直接用Apache Beam的S3 IO模块就能实现读取操作,不需要额外适配MinIO专属SDK,具体步骤和代码示例如下:

1. 添加必要的Maven依赖

在你的pom.xml中加入Beam的AWS IO依赖(版本请和你使用的Beam主版本保持一致,这里以2.54.0为例):

<dependency>
    <groupId>org.apache.beam</groupId>
    <artifactId>beam-sdks-java-io-amazon-web-services</artifactId>
    <version>2.54.0</version>
</dependency>

2. 完整代码示例

下面是可直接运行的代码,包含MinIO配置、Pipeline构建和文件读取逻辑,已经适配了play.min.io的测试环境:

import org.apache.beam.sdk.Pipeline;
import org.apache.beam.sdk.io.TextIO;
import org.apache.beam.sdk.options.PipelineOptions;
import org.apache.beam.sdk.options.PipelineOptionsFactory;
import org.apache.beam.sdk.values.PCollection;

public class MinIOBeamFileReader {
    public static void main(String[] args) {
        // 1. 初始化Pipeline配置,注入MinIO相关参数
        PipelineOptions options = PipelineOptionsFactory.create();
        
        // 配置play.min.io的连接信息,自己的MinIO实例直接替换对应参数即可
        options.as(org.apache.beam.sdk.io.aws.options.S3Options.class)
               .setS3EndpointOverride("https://play.min.io")
               .setAwsAccessKeyId("Q3AM3UQ867SPQQA43P2F")
               .setAwsSecretKey("zuf+tfteSlswRu7BJ86wekitnifILbZam1KYY3TG")
               .setPathStyleAccessEnabled(true); // MinIO默认需要开启路径风格访问

        // 2. 创建Pipeline实例
        Pipeline pipeline = Pipeline.create(options);

        // 3. 读取MinIO桶中的文件,替换成你自己的桶名,支持通配符匹配多个文件
        PCollection<String> fileContent = pipeline.apply(
            TextIO.read().from("s3://你的存储桶名称/*.txt")
        );

        // 4. 示例:打印读取到的文件内容,你可以替换成自己的业务处理逻辑
        fileContent.apply("Print File Content", org.apache.beam.sdk.transforms.ParDo.of(
            new org.apache.beam.sdk.transforms.DoFn<String, Void>() {
                @ProcessElement
                public void processElement(@Element String line) {
                    System.out.println("读取到内容:" + line);
                }
            }
        ));

        // 5. 启动Pipeline并等待执行完成
        pipeline.run().waitUntilFinish();
    }
}

关键配置说明

  • setS3EndpointOverride:指定MinIO的服务地址,play.min.io的公开地址是https://play.min.io,内网部署的MinIO替换成对应地址即可
  • setAwsAccessKeyId/setAwsSecretKey:play.min.io提供的公开测试密钥,自有MinIO实例替换成你自己的密钥
  • setPathStyleAccessEnabled(true):MinIO默认使用路径风格的存储地址(比如https://play.min.io/桶名/文件名),而AWS默认是虚拟主机风格,必须开启这个参数才能正常访问MinIO
  • TextIO.read().from("s3://你的存储桶名称/*.txt"):路径格式和S3完全一致,支持*通配符匹配多个文件,也可以指定单个文件路径比如s3://你的存储桶名称/one.txt

小提示

  • 生产环境中不要硬编码密钥,建议通过环境变量或者配置文件注入
  • 如果遇到网络连接问题,先确认是否能正常访问play.min.io(或者你的MinIO实例)
  • 不同Beam版本的参数名称可能略有差异,若报错可以查看对应版本的官方文档调整

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.07 14:08:15