Apache Beam TextIO Writer无法将SQS无界源数据写入目标位置
问题排查及解决方案
- 窗口触发策略缺失
固定窗口默认仅在水位线超过窗口结束时间后才会触发写入操作,你当前未配置触发策略,流运行状态下水位线推进不及时会导致窗口迟迟不闭合,临时文件不会被归档到目标路径。你可以给窗口配置显式的触发规则,示例如下:
.apply("Create Window",Window.into(FixedWindows.of(Duration.standardSeconds(10))) // 显式配置触发 .triggering(AfterWatermark.pastEndOfWindow() // 允许窗口结束前每1秒提前触发一次 .withEarlyFirings(AfterProcessingTime.pastFirstElementInPane() .plusDelayOf(Duration.standardSeconds(1)))) // 不允许数据延迟 .withAllowedLateness(Duration.ZERO) // 触发后丢弃当前pane数据 .discardingFiredPanes()
- 路径配置异常
检查getDestinationBucketUrl返回的GCS路径是否符合gs://桶名/路径/格式,临时目录和最终输出路径是否具有写入权限,建议将临时目录单独设置在目标路径的子目录下,避免跨目录移动文件无权限的问题:
.withTempDirectory(FileBasedSink.convertToFileResourceIfPossible(options.getDestinationBucketUrl() + "/temp"))
- 转换逻辑校验
确认SqsMessageToJson处理后输出结果非空,可通过RowPrinter的控制台打印内容确认转换后的消息是否正常生成,空数据不会触发写出文件。 - Shard配置优化
将withNumShards参数设置为1测试,或者设置为0由Beam自动分配shard数量,避免设置的shard数超过窗口内消息数量时,空shard不会生成输出文件。 - 运行模式配置
显式声明流运行模式,确保Direct Runner按流模式处理数据:
options.setRunner(DirectRunner.class); options.setStreaming(true);
修改后完整可测试代码:
Options options = PipelineOptionsFactory.fromArgs(CONFIG_STREAMING_SQS_GCS).withValidation().as(Options.class); // 显式指定运行模式 options.setRunner(DirectRunner.class); options.setStreaming(true); BasicAWSCredentials basicAWSCredentials = new BasicAWSCredentials("you-key", "your-secret"); options.setAwsCredentialsProvider(new AWSStaticCredentialsProvider(basicAWSCredentials)); Pipeline pipeline = Pipeline.create(options); pipeline.apply("Read messages from Sqs", SqsIO.read().withQueueUrl(options.getInputQueueUrl())) .apply("Get message contents", ParDo.of(new SqsMessageToJson())) .apply("Print incoming", ParDo.of(new RowPrinter<>("Print incoming"))) // 配置带触发规则的窗口 .apply("Create Window",Window.into(FixedWindows.of(Duration.standardSeconds(10))) .triggering(AfterWatermark.pastEndOfWindow() .withEarlyFirings(AfterProcessingTime.pastFirstElementInPane() .plusDelayOf(Duration.standardSeconds(1)))) .withAllowedLateness(Duration.ZERO) .discardingFiredPanes()) .apply("Write to GCS", TextIO.write() .withWindowedWrites() // 单独指定临时目录 .withTempDirectory(FileBasedSink.convertToFileResourceIfPossible(options.getDestinationBucketUrl() + "/temp")) .to(new WindowedFilenamePolicy(options.getOutputFilenamePrefix(), options.getShardTemplate(), options.getOutputFilenameSuffix()) .withSubDirectoryPolicy(options.getSubDirectoryPolicy())) // 测试阶段固定1个shard方便排查问题 .withNumShards(1)); PipelineResult run = pipeline.run(); run.waitUntilFinish();
内容的提问来源于stack exchange,提问作者Balasubramanian Naagarajan
相关产品推荐
相关产品推荐

