【问题标题】:Flink Streaming Compression not working using Amazon AWS S3 connector (StreamingFileSink with CompressWriterFactory)Flink 流压缩无法使用 Amazon AWS S3 连接器(StreamingFileSink 与 CompressWriterFactory)
【发布时间】:2020-09-20 03:10:46
【问题描述】:

当我将 Apache Flink 流式传输到 AWS S3 作为 Sink 时,标准版本 (forRowFormat) 工作正常。

StreamingFileSink<String> s3sink = StreamingFileSink
        .forRowFormat(new Path(s3Url),
                (String element, OutputStream stream) -> {
                    PrintStream out = new PrintStream(stream);
                    out.println(element);
                })
            .withBucketAssigner(new BucketAssigner())
            .withRollingPolicy(DefaultRollingPolicy.builder()
                    .withMaxPartSize(100)
                    .withRolloverInterval(30000)
                    .build())
            .withBucketCheckInterval(100)
            .build();

当我使用批量格式和 CompressWriterFactory 运行相同的东西时

StreamingFileSink<String> s3sink = StreamingFileSink
        .forBulkFormat(new Path(s3Url), 
                new CompressWriterFactory(new DefaultExtractor()))
        .withOutputFileConfig(outputFileConfig)
        .build();

它给了我以下错误。

(注意 - CompressWriterFactory 适用于 HDFS 方案“hdfs://host:port/path”)

java.lang.UnsupportedOperationException: S3RecoverableFsDataOutputStream cannot sync state to S3. Use persist() to create a persistent recoverable intermediate point.
    at org.apache.flink.fs.s3.common.utils.RefCountedBufferingFileStream.sync(RefCountedBufferingFileStream.java:112)
    at org.apache.flink.fs.s3.common.writer.S3RecoverableFsDataOutputStream.sync(S3RecoverableFsDataOutputStream.java:126)
    at org.apache.flink.formats.compress.writers.NoCompressionBulkWriter.finish(NoCompressionBulkWriter.java:56)
    at org.apache.flink.streaming.api.functions.sink.filesystem.BulkPartWriter.closeForCommit(BulkPartWriter.java:62)
    at org.apache.flink.streaming.api.functions.sink.filesystem.Bucket.closePartFile(Bucket.java:239)
    at org.apache.flink.streaming.api.functions.sink.filesystem.Bucket.prepareBucketForCheckpointing(Bucket.java:280)
    at org.apache.flink.streaming.api.functions.sink.filesystem.Bucket.onReceptionOfCheckpoint(Bucket.java:253)
    at org.apache.flink.streaming.api.functions.sink.filesystem.Buckets.snapshotActiveBuckets(Buckets.java:250)
    at org.apache.flink.streaming.api.functions.sink.filesystem.Buckets.snapshotState(Buckets.java:241)
    at org.apache.flink.streaming.api.functions.sink.filesystem.StreamingFileSink.snapshotState(StreamingFileSink.java:422)
    at org.apache.flink.streaming.util.functions.StreamingFunctionUtils.trySnapshotFunctionState(StreamingFunctionUtils.java:118)
    at org.apache.flink.streaming.util.functions.StreamingFunctionUtils.snapshotFunctionState(StreamingFunctionUtils.java:99)
    ...
    at org.apache.flink.streaming.runtime.tasks.StreamTask.runMailboxLoop(StreamTask.java:487)
    at org.apache.flink.streaming.runtime.tasks.StreamTask.invoke(StreamTask.java:470)
    at org.apache.flink.runtime.taskmanager.Task.doRun(Task.java:707)
    at org.apache.flink.runtime.taskmanager.Task.run(Task.java:532)

注意事项 -

  1. Flink 版本 1.10.0
  2. s3Url = "s3a://bucket/folder/path";

【问题讨论】:

    标签: amazon-s3 apache-flink flink-streaming


    【解决方案1】:

    这似乎是一个错误。您可以使用您在此处包含的描述打开 JIRA。

    【讨论】:

      【解决方案2】:

      刚刚发现了解决此问题的方法:指定 Hadoop 压缩编解码器:

      CompressWriters.forExtractor(new DefaultExtractor()).withHadoopCompression("GzipCodec")
      

      【讨论】:

        猜你喜欢
        • 2020-05-06
        • 2022-09-28
        • 1970-01-01
        • 2019-12-16
        • 1970-01-01
        • 2012-07-10
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多