【问题标题】:Spring Integration: Transform and route with headerSpring Integration:使用标头进行转换和路由
【发布时间】:2018-08-14 15:15:46
【问题描述】:

我正在构建一个基于 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;
  }

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

谢谢

【问题讨论】:

    标签: routing spring-integration spring-integration-dsl


    【解决方案1】:

    我相信您担心在 IssuanceTransformer 之后不再有所需的 AwsHeaders.STREAM 标头。当您开发自定义转换器时,您需要确保将请求消息中的所有标头传输到回复消息:与许多其他组件不同,转换器不会修改来自 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);
            }
    

    【讨论】:

      猜你喜欢
      • 2013-11-28
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2013-11-13
      • 2017-12-23
      相关资源
      最近更新 更多