【发布时间】:2022-01-22 16:06:12
【问题描述】:
我关注 (ZIP compressed input for Apache Flink) 并编写了以下代码片段,以使用简单的 TextInputFormat 处理目录中的 .gz 日志文件。它适用于我的本地测试目录,扫描并自动打开.gz 文件内容。但是,当我使用 s3 存储桶源运行它时,它不会处理 .gz 压缩文件。不过,这个 Flink 作业仍然会打开 s3 存储桶上的 .log 文件。似乎它只是不解压缩 .gz 文件。如何在 s3 文件系统上解决此问题?
public static void main(String[] args) throws Exception {
final ParameterTool params = ParameterTool.fromArgs(args);
final String sourceLogDirPath = params.get("source_log_dir_path", "s3://my-test-bucket-logs/"); // "/Users/my.user/logtest/logs"
final Long checkpointInterval = Long.parseLong(params.get("checkpoint_interval", "60000"));
final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime);
env.enableCheckpointing(checkpointInterval, CheckpointingMode.EXACTLY_ONCE);
env.getCheckpointConfig().enableExternalizedCheckpoints(
CheckpointConfig.ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION);
env.getConfig().setGlobalJobParameters(params);
TextInputFormat textInputFormat = new TextInputFormat(new Path(sourceLogDirPath));
textInputFormat.setNestedFileEnumeration(true);
DataStream<String> stream = env.readFile(
textInputFormat, sourceLogDirPath,
FileProcessingMode.PROCESS_CONTINUOUSLY, 100);
stream.print();
env.execute();
}
这是我的类路径 jar flink 库:
/opt/flink/lib/flink-csv-1.13.2.jar:/opt/flink/lib/flink-json-1.13.2.jar:/opt/flink/lib/flink-shaded-zookeeper- 3.4.14.jar:/opt/flink/lib/flink-table-blink_2.12-1.13.2.jar:/opt/flink/lib/flink-table_2.12-1.13.2.jar:/opt/flink /lib/log4j-1.2-api-2.12.1.jar:/opt/flink/lib/log4j-api-2.12.1.jar:/opt/flink/lib/log4j-core-2.12.1.jar:/ opt/flink/lib/log4j-slf4j-impl-2.12.1.jar:/opt/flink/lib/sentry_log4j2_deploy.jar:/opt/flink/lib/flink-dist_2.12-1.13.2.jar:::
附:我也尝试了s3a://<bucket>/,但没有成功。
【问题讨论】:
标签: amazon-s3 apache-flink flink-streaming