代码之家  ›  专栏  ›  技术社区  ›  sansari

Spring集成:带标题的转换和路由

  •  2
  • sansari  · 技术社区  · 8 年前

    我正在构建一个基于Spring的库,在将消息转换为正确的类型之后,它应该使用消息并将其传递到配置的通道。我的库可以通过“streamToConsume:finalChannelDestination”对列表进行配置。

    streams:
        source: destinationChannel
    

    我想要一个 IntegrationFlow 如下所示:

    IntegrationFlows
            .from(kinesisInboundChannelAdapter(amazonKinesis(), streamNames))
            .transform(new IssuanceTransformer())
            .route(router())
            .get();
    
    public HeaderValueRouter router() {
        HeaderValueRouter router = new HeaderValueRouter(AwsHeaders.STREAM);
        consumerClientProperties.getKinesis().getStreams().forEach((k, v) ->
            router.setChannelMapping(k, v)
        );
        return router;
      }
    

    转换事件,然后将其传递到配置中映射到流的通道。如何在转换后保留事件头以便将其发送到正确的通道?

    谢谢你

    1 回复  |  直到 8 年前
        1
  •  1
  •   Artem Bilan    8 年前

    我相信你的担心 IssuanceTransformer 没有一个愿望 AwsHeaders.STREAM 标题。开发自定义转换器时,需要确保将所有头从请求消息传输到回复消息:与许多其他组件不同,Transformer不会修改来自POJO的回复消息。

    为此,您可以使用如下内容:

    MessageBuilder.withPayload(myPayload).copyHeadersIfAbsent(requestMessage.getHeaders()).build();
    

    注意:您可以使用 AwsHeaders.RECEIVED_STREAM 因为这个是从 KinesisMessageDrivenChannelAdapter :

    private void performSend(AbstractIntegrationMessageBuilder<?> messageBuilder, Object rawRecord) {
            messageBuilder.setHeader(AwsHeaders.RECEIVED_STREAM, this.shardOffset.getStream())
                    .setHeader(AwsHeaders.SHARD, this.shardOffset.getShard());
    
            if (CheckpointMode.manual.equals(KinesisMessageDrivenChannelAdapter.this.checkpointMode)) {
                messageBuilder.setHeader(AwsHeaders.CHECKPOINTER, this.checkpointer);
            }