【问题标题】:spring integration - how to aggregate split lines into batch of x?弹簧集成 - 如何将分割线聚合成一批 x?
【发布时间】:2018-04-19 09:15:03
【问题描述】:

我们需要将事件发送到 kinesis,由于 AWS 定价,我们计划将记录分批放入 kinesis。

我们读入一个 csv 文件,然后使用文件拆分器将行吐出并将每一行转换为 json。

那么在转换为 json 之后,我们如何将这些行批处理为每批 25 行,以便我们的 kinesis serviceActivator 可以发送批处理?

任何例子都将不胜感激。

    <int-file:splitter id="fileLineSplitter"
                       input-channel="fileInputChannel"
                       output-channel="splitterOutputChannel"
                       markers="true" />


<int:transformer id="csvToDataCdrTransformer"
                     ref="dataCdrLineTransformer"
                     method="transform"
                     input-channel="lineOutputChannel"
                     output-channel="dataCdrObjectInputChannel">
    </int:transformer>


    <int:object-to-json-transformer input-channel="dataCdrObjectInputChannel"
                                    output-channel="kinesisSendChannel">
        <int:poller fixed-delay="50"/>
    </int:object-to-json-transformer>

编辑:我按照“Artem Bilan”的建议添加了它,它起作用了

<int:aggregator input-channel="aggregateChannel"
                output-channel="toJsonChannel"
                release-strategy-expression="#this.size() eq 2"
                expire-groups-upon-completion="true"/>

但我得到错误:

  1. 我正在使用 markers="true",以便我们知道它的文件结尾,因此我们可以将其重命名为“.done”。

  2. 在拆分器和转换器之间添加了一个路由器,当 FileMarker 为 END 时,它仅路由到“nullChannel”或“fileProcessedChannel”,否则,拆分线进入 default-output-channel="lineOutputChannel"

    <int:router ref="fileMarkerCustomRouter" inputchannel="splitterOutputChannel" default-output-channel="lineOutputChannel"/>
    

路由器代码如下所示

 @Override
    protected Collection<MessageChannel> determineTargetChannels(Message<?> message) {
        Collection<MessageChannel> targetChannels = new ArrayList<MessageChannel>();

        if (isPayloadTypeFileMarker(message)) {

            FileSplitter.FileMarker payload = (FileSplitter.FileMarker) message.getPayload();

            if (isStartOfFile(payload)) {

                targetChannels.add(nullChannel);

            } else if (isEndOfFile(payload)) {

                targetChannels.add(fileProcessedChannel);
            }
        }
        return targetChannels;
    }

但出现此错误:

Caused by: java.lang.IllegalStateException: Null correlation not allowed.  Maybe the CorrelationStrategy is failing?

有什么想法吗?

【问题讨论】:

    标签: spring-boot spring-integration spring-integration-aws


    【解决方案1】:

    为此,您绝对需要一个&lt;aggregator&gt;release-strategy-expression="25"expire-groups-upon-completion="true",以便在发布一个correlationKey 后为其组成一个新组。

    不,确定您为什么需要markers="true",但没有&lt;int-file:splitter&gt; 会填充适当的相关标头。因此,您甚至可以考虑在之后仅依赖默认拆分和默认聚合。

    此外,您应该考虑将来自聚合器的结果转换为 JSON。它发出一个List&lt;?&gt;。将整个列表序列化为 JSON 非常有效。另外,在发送到 Kinesis 之前,您可能需要再进行一次转换。

    因此,您的配置原型应该是这样的:

    <int-file:splitter id="fileLineSplitter"
                       input-channel="fileInputChannel"
                       output-channel="splitterOutputChannel"/>
    
    <int:transformer id="csvToDataCdrTransformer"
                     ref="dataCdrLineTransformer"
                     method="transform"
                     input-channel="lineOutputChannel"
                     output-channel="aggregateChannel">
    </int:transformer>
    
    <int:aggregator input-channel="aggregateChannel" 
                    output-channel="toJsonChannel"
                    expire-groups-upon-completion="true" />
    
    <int:object-to-json-transformer input-channel="toJsonChannel"
                                    output-channel="kinesisSendChannel"/>
    

    这样整个文件将被视为一个批处理。您拆分它,处理每一行,将它们聚合回列表,然后在发送到 Kinesis 之前转换为 JSON。

    从这里我想请你提出一个 JIRA 来添加ObjectToJsonTransformer.ResultType.BYTES 模式,以便更好地使用基于byte[] 的下游组件,例如KinesisMessageHandler

    【讨论】:

    • 我正在使用 marker="true" 以便我们知道它是文件的结尾,因此我们可以将其重命名为“.done”。
    • 我应用了您的解决方案,它确实有效。但是,我不得不在分离器和变压器之间放置一个路由器,并且出现错误,请您帮我解决我对这个问题所做的“编辑”。
    • 如果您的transform 方法返回Message&lt;?&gt;,您负责将输入标头复制到输出消息。如果您只返回新的有效负载,框架将复制标头 - 拆分器默认添加 correlationId 标头。
    • 我添加了一个路由器,请参阅上面编辑部分 2 中的路由器代码,并收到“不允许空关联”错误。上面的方法 determineTargetChannels(...) 从 Splitter 接收带有相关 ID 的 Message> 并且当路由器返回一个通道时,我如何将相关 ID 作为返回通道的一部分传回?我正在使用路由器,以便我知道它何时是文件的结尾,因此标记 = true,以便我可以重命名文件。你会建议另一种重命名文件的方法吗?此外,如果我卸下路由器并从 Splitter 转到 Trnasformer,那么一切都按照@Artem 的建议进行。
    • 使用markers = true,您必须在FileSplitter 上手动打开correlationapply-sequence="true"
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-04-08
    • 1970-01-01
    • 2014-11-10
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多