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

使用TextIO.write()向S3写入数据时遇AbstractMethodError求助

问题:使用TextIO.write()写入S3时重命名临时文件失败

我编写了如下Google Dataflow代码,尝试将内容写入S3:

import org.apache.beam.sdk.Pipeline;
import org.apache.beam.sdk.io.FileSystems;
import org.apache.beam.sdk.io.TextIO;
import org.apache.beam.sdk.io.fs.ResourceId;
import org.apache.beam.sdk.options.*;
import org.apache.beam.sdk.transforms.Create;

public class S3Test {

    public static void main(String[] args) {
        PipelineOptions options = PipelineOptionsFactory.fromArgs(args).withValidation().create();
        Pipeline pipeline = Pipeline.create(options);

        ResourceId outputDir = FileSystems.matchNewResource("s3://my-bucket/temp/", true);

        pipeline
                .apply("New",
                        Create.of("Hello World!!")
                )
                .apply(
                        "Write to S3",
                        TextIO.write()
                                .to("s3://my-bucket/test.txt")
                                .withoutSharding()
                                .withTempDirectory(outputDir)
                );

        pipeline.run();
    }
}

代码能在S3桶中创建临时文件,但在将临时文件重命名为最终目标文件名时,抛出如下错误:

Receiver class org.apache.beam.sdk.io.aws.s3.S3FileSystem does not define or inherit an implementation of the resolved method 'abstract void rename(java.util.List, java.util.List, org.apache.beam.sdk.io.fs.MoveOptions[])' of abstract class org.apache.beam.sdk.io.FileSystem.

我猜测需要重写或实现rename()方法,但不知从何入手。有没有人成功使用TextIO.write()写入S3?

附完整运行日志:

Dec 24, 2022 4:56:34 PM
org.apache.beam.sdk.io.WriteFiles$WriteShardsIntoTempFilesFn
processElement INFO: Opening writer
8a0dbf70-a0fc-48dd-92f9-8ff4000a8076 for window
org.apache.beam.sdk.transforms.windowing.GlobalWindow@2dddc1b9 pane
PaneInfo{isFirst=true, isLast=true, timing=ON_TIME, index=0,
onTimeIndex=0} destination null Dec 24, 2022 4:56:37 PM
org.apache.beam.sdk.io.FileBasedSink$Writer close INFO: Successfully
wrote temporary file
s3://my-bucket/temp/.temp-beam-13da7e3e-15a9-4285-a90c-30c7c6dcfe99/9e53ec6a8a0dbf70-a0fc-48dd-92f9-8ff4000a8076
Dec 24, 2022 4:56:37 PM
org.apache.beam.sdk.io.WriteFiles$FinalizeTempFileBundles$FinalizeFn
process INFO: Finalizing 1 file results Dec 24, 2022 4:56:37 PM
org.apache.beam.sdk.io.FileBasedSink$WriteOperation
createMissingEmptyShards INFO: Finalizing for destination null num
shards 1. Dec 24, 2022 4:56:37 PM
org.apache.beam.sdk.io.FileBasedSink$WriteOperation moveToOutputFiles
INFO: Will copy temporary file
FileResult{tempFilename=s3://my-bucket/temp/.temp-beam-13da7e3e-15a9-4285-a90c-30c7c6dcfe99/9e53ec6a8a0dbf70-a0fc-48dd-92f9-8ff4000a8076,
shard=0,
window=org.apache.beam.sdk.transforms.windowing.GlobalWindow@2dddc1b9,
paneInfo=PaneInfo{isFirst=true, isLast=true, timing=ON_TIME, index=0,
onTimeIndex=0}} to final location s3://my-bucket/test.txt Exception in
thread "main" org.apache.beam.sdk.Pipeline$PipelineExecutionException:
java.lang.AbstractMethodError: Receiver class
org.apache.beam.sdk.io.aws.s3.S3FileSystem does not define or inherit
an implementation of the resolved method 'abstract void
rename(java.util.List, java.util.List,
org.apache.beam.sdk.io.fs.MoveOptions[])' of abstract class
org.apache.beam.sdk.io.FileSystem. at
org.apache.beam.runners.direct.DirectRunner$DirectPipelineResult.waitUntilFinish(DirectRunner.java:374)
at
org.apache.beam.runners.direct.DirectRunner$DirectPipelineResult.waitUntilFinish(DirectRunner.java:342)
at
org.apache.beam.runners.direct.DirectRunner.run(DirectRunner.java:218)
at
org.apache.beam.runners.direct.DirectRunner.run(DirectRunner.java:67)
at org.apache.beam.sdk.Pipeline.run(Pipeline.java:323) at
org.apache.beam.sdk.Pipeline.run(Pipeline.java:309) at
com.example.S3Test.main(S3Test.java:255) Caused by:
java.lang.AbstractMethodError: Receiver class
org.apache.beam.sdk.io.aws.s3.S3FileSystem does not define or inherit
an implementation of the resolved method 'abstract void
rename(java.util.List, java.util.List,
org.apache.beam.sdk.io.fs.MoveOptions[])' of abstract class
org.apache.beam.sdk.io.FileSystem. at
org.apache.beam.sdk.io.FileSystems.renameInternal(FileSystems.java:323)
at org.apache.beam.sdk.io.FileSystems.rename(FileSystems.java:308)
at
org.apache.beam.sdk.io.FileBasedSink$WriteOperation.moveToOutputFiles(FileBasedSink.java:802)
at
org.apache.beam.sdk.io.WriteFiles$FinalizeTempFileBundles$FinalizeFn.process(WriteFiles.java:1077)


解决方案

核心原因

这个错误并非需要你自己重写rename方法,而是Beam核心SDK与AWS S3扩展SDK版本不兼容导致的:Beam的FileSystem抽象类在某个版本新增了带MoveOptions[]参数的rename方法,但你使用的S3FileSystem实现类(来自beam-sdks-java-io-amazon-web-services)版本较低,未实现该新方法,从而触发AbstractMethodError。

具体修复步骤

  1. 对齐依赖版本
    确保所有Beam相关依赖的版本完全一致,包括核心SDK、AWS S3扩展、运行器(如DirectRunner)。例如统一使用2.40.0版本:
<!-- Maven依赖示例 -->
<dependencies>
    <!-- Beam核心SDK -->
    <dependency>
        <groupId>org.apache.beam</groupId>
        <artifactId>beam-sdks-java-core</artifactId>
        <version>2.40.0</version>
    </dependency>
    <!-- Beam AWS S3 IO扩展 -->
    <dependency>
        <groupId>org.apache.beam</groupId>
        <artifactId>beam-sdks-java-io-amazon-web-services</artifactId>
        <version>2.40.0</version>
    </dependency>
    <!-- 本地运行用DirectRunner -->
    <dependency>
        <groupId>org.apache.beam</groupId>
        <artifactId>beam-runners-direct-java</artifactId>
        <version>2.40.0</version>
        <scope>runtime</scope>
    </dependency>
</dependencies>
  1. 验证凭证配置
    确保AWS凭证已正确配置(如通过环境变量AWS_ACCESS_KEY_ID/AWS_SECRET_ACCESS_KEY,或~/.aws/credentials文件),避免后续出现权限类错误。

  2. 重新运行代码
    版本对齐后,你的原有代码逻辑无需修改,重新运行即可成功完成临时文件到最终文件的重命名操作。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 07:15:35