【问题标题】:Camel: How to split and then aggregate when number of item is less than batch sizeCamel:当项目数小于批量大小时如何拆分然后聚合
【发布时间】:2019-05-29 18:35:33
【问题描述】:

我有一个从 S3 读取文件的 Camel 路由,并按如下方式处理输入文件:

  1. 使用Bindy 将每一行解析为一个 POJO(学生)
  2. 通过body()分割输出
  3. 按正文的属性 (.semester) 和 2 的批量大小
  4. 调用持久化服务以给定批次上传到数据库

问题在于,批量大小为 2 且记录数为奇数时,总会有一条记录未保存。

提供的代码是 Kotlin,但不应与等效的 Java 代码有很大不同(在 "\${simple expression}" 前面加上斜线或缺少分号来终止语句。

如果我将批量大小设置为 1,则保存每条记录,否则永远不会保存最后一条记录。

我已经检查了几次message-processor 的文档,但它似乎没有涵盖这种特殊情况。

除了completionSize之外,我还设置了[completionTimeout|completionInterval],但没有任何区别。

以前有人遇到过这个问题吗?

val csvDataFormat = BindyCsvDataFormat(Student::class.java)

from("aws-s3://$student-12-bucket?amazonS3Client=#amazonS3&delay=5000")
    .log("A new Student input file has been received in S3: '\${header.CamelAwsS3BucketName}/\${header.CamelAwsS3Key}'")
    .to("direct:move-input-s3-object-to-in-progress")
    .to("direct:process-s3-file")
    .to("direct:move-input-s3-object-to-completed")
    .end()

from("direct:process-s3-file")
    .unmarshal(csvDataFormat)
    .split(body())
    .streaming()
    .parallelProcessing()
    .aggregate(simple("\${body.semester}"), GroupedBodyAggregationStrategy())
    .completionSize(2)
    .bean(persistenceService)
    .end()

输入 CSV 文件包含七 (7) 条记录,这是生成的输出(添加了一些调试日志记录):

WARN 19540 --- [student-12-move] c.a.s.s.internal.S3AbortableInputStream :并非所有字节都从 S3ObjectInputStream 中读取,因此中止 HTTP 连接。这很可能是一个错误,并可能导致次优行为。仅通过远程 GET 请求您需要的字节,或在使用后排空输入流。 INFO 19540 --- [student-12-move] student-workflow-main:在 S3 中收到了一个新的学生输入文件:'student-12-bucket/inbox/foo.csv' INFO 19540 --- [student-12-move] move-input-s3-object-to-in-progress:将 S3 文件“inbox/foo.csv”移动到“in-progress”文件夹... INFO 19540 --- [student-12-move] student-workflow-main:将输入 S3 文件“in-progress/foo.csv”移动到“in-progress”文件夹... INFO 19540 --- [student-12-move] pre-process-s3-file-records:开始保存到数据库... 调试 19540 --- [读取 #7 - 拆分] c.b.i.d.s.StudentPersistenceServiceImpl :将记录保存到数据库:学生(id=7,姓名=学生 7,学期=2,javaMarks=25) 调试 19540 --- [读取 #7 - 拆分] c.b.i.d.s.StudentPersistenceServiceImpl :将记录保存到数据库:学生(id=5,姓名=学生 5,学期=2,javaMarks=81) 调试 19540 --- [读取 #3 - 拆分] c.b.i.d.s.StudentPersistenceServiceImpl :将记录保存到数据库:学生(id=6,姓名=学生 6,学期=1,javaMarks=15) 调试 19540 --- [读取 #3 - 拆分] c.b.i.d.s.StudentPersistenceServiceImpl:将记录保存到数据库:学生(id=2,姓名=学生 2,学期=1,javaMarks=62) 调试 19540 --- [读取 #2 - 拆分] c.b.i.d.s.StudentPersistenceServiceImpl:将记录保存到数据库:学生(id=3,姓名=学生 3,学期=2,javaMarks=72) 调试 19540 --- [读取 #2 - 拆分] c.b.i.d.s.StudentPersistenceServiceImpl:将记录保存到数据库:学生(id=1,姓名=学生 1,学期=2,javaMarks=87) INFO 19540 --- [student-12-move] device-group-workflow-main:结束预处理 S3 CSV 文件记录... INFO 19540 --- [student-12-move] move-input-s3-object-to-completed:将 S3 文件“进行中/foo.csv”移动到“已完成”文件夹... INFO 19540 --- [student-12-move] device-group-workflow-main:将 S3 文件“in-progress/foo.csv”移动到“已完成”文件夹...

【问题讨论】:

  • completionTimeout 应该在超时时触发最后一行。如果这不起作用,那就奇怪了。
  • 如果我将 simple("${body.semester}") 替换为 constant(true),它确实表现出正确的行为。这可能是一个错误......
  • 您使用的是什么版本的骆驼? group key 是否 body.semester 或 constant 不应该影响超时。
  • 我正在使用以下组件:+ camel.version = 2.23.0 + spring-boot.version = 2.1.1.RELEASE + kotlin.version = 1.3.10 + aws-java-sdk。版本 = 1.11.461
  • 并且您确定您的 bean 中没有什么问题仅适用于 1 条记录。您是否尝试过在聚合之后添加一个日志,以查看它在触发完成超时等时记录了一些内容。如果仍然存在问题,您可以尝试在 github 上构建一个项目,以便其他人更轻松地查看。并且可以在没有 AWS 帐户等的情况下轻松运行。

标签: kotlin apache-camel spring-camel


【解决方案1】:

如果您需要立即完成您的消息,那么您可以指定一个完成谓词,该谓词基于拆分器设置的交换属性。我没试过,但我认为

.completionPredicate( simple( "${exchangeProperty.CamelSplitComplete}" ) )

将处理最后一条消息。

我的另一个担心是您在拆分器中设置了parallelProcessing,这可能意味着消息没有按顺序处理。它真的是您希望应用并行处理的拆分器,还是实际上是聚合器?除了聚合它们,然后处理它们之外,您似乎没有对拆分记录做任何事情,因此将 parallelProcessing 指令移动到聚合器可能会更好。

【讨论】:

  • 这并没有解决问题,但它帮助我了解了问题的原因。问题似乎是,当按主体的某个属性分组时,最后可能会有一些结转交换,而这些交换完成条件都不会为真:*交换完成大小 == 2 * exchangeProperty.CamelSplitComplete ==真的不知道如何解决它..
  • 可能关闭并行处理,因为您可以在启用它的情况下进行乱序处理。
  • 我确实按照建议在拆分期间关闭了并行处理。最终结果是一样的..
猜你喜欢
  • 2023-03-26
  • 1970-01-01
  • 1970-01-01
  • 2020-11-04
  • 1970-01-01
  • 2020-04-25
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多