Flink测试作业执行报错No FileSystem for scheme "s3"求助
大家好,我在运行Flink测试作业时遇到了棘手的问题,一直提示No FileSystem for scheme "s3",折腾了半天没解决,来请教下各位!
先把完整的报错信息贴出来:
WARNING: An illegal reflective access operation has occurred
WARNING: Illegal reflective access by org.apache.flink.streaming.runtime.translators.DataStreamV2SinkTransformationTranslator (file:/opt/flink/lib/flink-dist-1.20.1.jar) to field java.util.Collections$UnmodifiableMap.m
WARNING: Please consider reporting this to the maintainers of org.apache.flink.streaming.runtime.translators.DataStreamV2SinkTransformationTranslator
WARNING: Use --illegal-access=warn to enable warnings of further illegal reflective access operations
WARNING: All illegal access operations will be denied in a future releaseThe program finished with the following exception:
org.apache.flink.client.program.ProgramInvocationException: The main method caused an error: No FileSystem for scheme "s3"
at org.apache.flink.client.program.PackagedProgram.callMainMethod(PackagedProgram.java:373)
at org.apache.flink.client.program.PackagedProgram.invokeInteractiveModeForExecution(PackagedProgram.java:223)
at org.apache.flink.client.ClientUtils.executeProgram(ClientUtils.java:113)
at org.apache.flink.client.cli.CliFrontend.executeProgram(CliFrontend.java:1026)
at org.apache.flink.client.cli.CliFrontend.run(CliFrontend.java:247)
at org.apache.flink.client.cli.CliFrontend.parseAndRun(CliFrontend.java:1270)
at org.apache.flink.client.cli.CliFrontend.lambda$mainInternal$10(CliFrontend.java:1367)
at org.apache.flink.runtime.security.contexts.NoOpSecurityContext.runSecured(NoOpSecurityContext.java:28)
at org.apache.flink.client.cli.CliFrontend.mainInternal(CliFrontend.java:1367)
at org.apache.flink.client.cli.CliFrontend.main(CliFrontend.java:1335)
Caused by: org.apache.hadoop.fs.UnsupportedFileSystemException: No FileSystem for scheme "s3"
at org.apache.hadoop.fs.FileSystem.getFileSystemClass(FileSystem.java:3443)
at org.apache.hadoop.fs.FileSystem.createFileSystem(FileSystem.java:3466)
at org.apache.hadoop.fs.FileSystem.access$300(FileSystem.java:174)
at org.apache.hadoop.fs.FileSystem$Cache.getInternal(FileSystem.java:3574)
at org.apache.hadoop.fs.FileSystem$Cache.get(FileSystem.java:3521)
at org.apache.hadoop.fs.FileSystem.get(FileSystem.java:540)
at org.apache.hadoop.fs.Path.getFileSystem(Path.java:365)
at org.apache.parquet.hadoop.ParquetReader$Builder.build(ParquetReader.java:391)
at com.xyz.odl.util.parquet.JsonParquetReader.readParquetFileAsCollection(JsonParquetReader.java:53)
at com.xyz.odl.BackfillDynamoJob.main(BackfillDynamoJob.java:47)
at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke(Unknown Source)
at java.base/jdk.internal.reflect.DelegatingMethodAccessorImpl.invoke(Unknown Source)
at java.base/java.lang.reflect.Method.invoke(Unknown Source)
at org.apache.flink.client.program.PackagedProgram.callMainMethod(PackagedProgram.java:356)
... 9 more
我的测试环境是用Testcontainers启动Flink JobManager容器的,下面是对应的测试代码:
private GenericContainer<?> constructAndStartFlinkJobManager() { String flinkConfYamlPath = Paths.get("src/test/resources/flink-conf.yaml").toAbsolutePath().toString(); String flinkS3FsPrestoJar = Paths.get("src/test/resources/plugins/flink-s3-fs-presto-1.20.1.jar") .toAbsolutePath() .toString(); try { GenericContainer<?> flinkJobManagerContainer = new GenericContainer<>(DockerImageName.parse("arm64v8/flink:1.20.1-java11")) .withCopyFileToContainer( MountableFile.forHostPath(propertiesJsonFile.getAbsolutePath()), "/opt/flink/conf/application-properties.json") .withCopyFileToContainer( MountableFile.forHostPath(flinkConfYamlPath), "/opt/flink/conf/flink-conf.yaml") .withNetwork(network) .withNetworkAliases("jobmanager") .withCommand("jobmanager") .withExposedPorts(8081) .withEnv("JOB_MANAGER_RPC_ADDRESS", "jobmanager") .withEnv( Map.of( "AWS_ACCESS_KEY_ID", "test-access-key", "AWS_SECRET_ACCESS_KEY", "test-secret-key", "AWS_REGION", "us-east-1", "AWS_ENDPOINT_URL_S3", "http://localstack:4566")) .waitingFor( Wait.forHttp("/").forPort(8081).withStartupTimeout(Duration.ofSeconds(120))); flinkJobManagerContainer.start(); ExecResult mkdirResult = flinkJobManagerContainer.execInContainer( "mkdir", "-p", "/opt/flink/plugins/s3-fs-presto"); if (mkdirResult.getExitCode() != 0) { throw new RuntimeException("Failed to create directory: " + mkdirResult.getStderr()); } ExecResult execResult = flinkJobManagerContainer.execInContainer( "cp", "/opt/flink/opt/flink-s3-fs-presto-1.20.1.jar", "/opt/flink/plugins/s3-fs-presto/"); if (execResult.getExitCode() != 0) { throw new RuntimeException("Failed to copy plugin JAR: " + execResult.getStderr()); } return flinkJobManagerContainer; } catch (IOException | InterruptedException | ContainerException e) { throw new RuntimeException("Failed to start Flink JobManager container", e); } }
我已经尝试把s3-fs-presto的jar包复制到容器的plugins/s3-fs-presto目录下,也配置了AWS相关的环境变量,但还是报这个错误。有没有大佬能帮我排查下哪里出问题了?
内容来源于stack exchange

