【问题标题】:Spring Batch : One Reader, composite processor (two classes with different entities) and two kafkaItemWriterSpring Batch:一个Reader,复合处理器(两个具有不同实体的类)和两个kafkaItemWriter
【发布时间】:2023-03-27 19:19:02
【问题描述】:

ItemReader 正在从 DB2 读取数据并提供 java 对象 ClaimDto。现在ClaimProcessor 接收ClaimDto 的对象并返回由claimRecord1claimRecord2 组成的CompositeClaimRecord 对象,该对象将被发送到两个不同的Kafka 主题。如何将claimRecord1claimRecord2分别写入topic1和topic2。

【问题讨论】:

  • 通过写一个ItemWriter来做到这一点。
  • 完全正确,但我无法弄清楚如何在编写器类中添加两个不同实体的委托,并从 CompositeClaimRecord 中获取每个实体以写入对应的主题。
  • 1,作家,2个主题,获取字段值,放入主题?为什么这么难?
  • 你能分享一些spring-batch应用的例子吗?我正在使用 KafkaItemWriter 类并尝试添加两个键相同但值不同的委托。
  • 不要使用KafkaItemWriter。或者在您自己编写的复合项目编写器中使用 2。

标签: spring-boot apache-kafka spring-batch kafka-producer-api itemwriter


【解决方案1】:

只需编写一个自定义的ItemWriter 即可。

public class YourItemWriter implements ItemWriter<CompositeClaimRecord>` {

  private final ItemWriter<Record1> writer1;
  private final ItemWriter<Record2> writer2;

  public YourItemWriter(ItemWriter<Record1> writer1, ItemWriter<Record2> writer2>) {
    this.writer1=writer1;
    this.writer2=writer2;
}

  public void write(List<CompositeClaimRecord> items) throws Exception {

    for (CompositeClaimRecord record : items) {
       writer1.write(Collections.singletonList(record.claimRecord1));
       writer2.write(Collections.singletonList(record.claimRecord2));

    }
  }
}

或者不是一次写入 1 条记录,而是将单个列表转换为 2 个列表并传递。但是这样处理错误可能有点挑战。 \

public class YourItemWriter implements ItemWriter<CompositeClaimRecord>` {

  private final ItemWriter<Record1> writer1;
  private final ItemWriter<Record2> writer2;

  public YourItemWriter(ItemWriter<Record1> writer1, ItemWriter<Record2> writer2>) {
    this.writer1=writer1;
    this.writer2=writer2;
}

  public void write(List<CompositeClaimRecord> items) throws Exception {

    List<ClaimRecord1> record1List = items.stream().map(it -> it.claimRecord1).collect(Collectors.toList());
    List<ClaimRecord2> record2List = items.stream().map(it -> it.claimRecord2).collect(Collectors.toList());

    writer1.write(record1List);
    writer2.write(record2List);


  }
}

【讨论】:

  • 你会在工作的主要步骤中直接调用 YourItemWriter 还是将委托添加到编写器,然后在 YourItemWriter 的 write 方法中使用它们?
  • 对不起,我提供了编写器的整个实现,还有什么不清楚的地方?
  • 我可以在您的代码之后写两个主题。如何在单个事务中编写它们?
【解决方案2】:

您可以使用ClassifierCompositeItemWriter 和两个KafkaItemWriters 作为代表(每个主题一个)。

Classifier 将根据项目的类型(claimRecord1claimRecord2)对项目进行分类,并将它们路由到相应的 kafka 项目编写器(topic1topic2)。

【讨论】:

  • 你能分享一个例子吗?我只能看到具有相同实体类型而不具有不同实体类型的文章。
  • 看起来我误解了你的问题,因为它不清楚。在标题中你说composite processor (two classes with different entities),但从描述来看,你似乎在CompositeClaimRecord 中封装了不同类型的东西。接受ClaimDto 并返回CompositeClaimRecord 的处理器是常规处理器,而不是复合处理器。 CompositeItemProcessor 接受委托处理器列表,并一次将项目传递给该处理器链。
  • 此答案不适用于您基于封装的方法,@M 的答案。 Deinum 是要走的路。我投了赞成票。
猜你喜欢
  • 2020-07-08
  • 2013-03-03
  • 1970-01-01
  • 1970-01-01
  • 2019-10-07
  • 1970-01-01
  • 2012-12-20
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多