使用TextIO.write()向S3写入数据时遇AbstractMethodError求助
我编写了如下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。
具体修复步骤
- 对齐依赖版本
确保所有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>
验证凭证配置
确保AWS凭证已正确配置(如通过环境变量AWS_ACCESS_KEY_ID/AWS_SECRET_ACCESS_KEY,或~/.aws/credentials文件),避免后续出现权限类错误。重新运行代码
版本对齐后,你的原有代码逻辑无需修改,重新运行即可成功完成临时文件到最终文件的重命名操作。
内容的提问来源于stack exchange,提问作者Shahid Thaika

