Spring Integration - Missing messages while routing to DirectChannel

Viewed 81

I have one DirectChannel(CHANNEL-ABC), Multiple producers put messages to it and there is only one consumer of that channel.

My Application polls one folder for files in it and CHANNEL-1 is called per file.

Producer as follow.

IntegrationFlows.from("CHANNEL-1")
                .handle(
                   ...
                   logger.info("CHANNEL-1 flow started")
                   ...
                 )
                .split()
                .handle(
                   ...
                 )
                .aggregate(...)
                .handle(
                   ...
                   logger.info("Aggregated")
                 )
                .channel("CHANNEL-ABC")
                .handle(
                   ...
                   logger.info("CHANNEL-1 flow ended")
                   ...
                 )
                .get();

IntegrationFlows.from("CHANNEL-2")
                .handle(
                   ...
                   logger.info("CHANNEL-2 flow started")
                   ...
                 )
                .handle(
                   ...
                 )
                .channel("CHANNEL-ABC")
                .handle(
                   ...
                   logger.info("CHANNEL-2 flow ended")
                   ...
                 )
                .get();

IntegrationFlows.from("CHANNEL-3")
                .handle(
                   ...
                   logger.info("CHANNEL-3 flow started")
                   ...
                 )
                .channel("CHANNEL-ABC")
                .handle(
                   ...
                   logger.info("CHANNEL-3 flow ended")
                   ...
                 )
                .get();

Note: All CHANNEL-1, CHANNEL-2 and CHANNEL-3 are DirectChannel.

Consumer as follow.

IntegrationFlows.from("CHANNEL-ABC")
                .handle(
                   ...
                   logger.info("CHANNEL-ABC flow started")
                   ...
                 )
                .handle(....)
                .handle(
                   ...
                   logger.info("CHANNEL-3 flow ended")
                   ...
                 )
                .get();

I had put multiple files in polling folder and I have observed logs as below.

For File-1 //CHANNEL-ABC called but route after channel calling not executed

CHANNEL-1 flow started
Aggregated
CHANNEL-ABC flow started
CHANNEL-ABC flow ended

For File-2 //CHANNEL-ABC flow not executed and routed to another rout after that

CHANNEL-1 flow started
Aggregated
CHANNEL-2 flow ended

For File-3 //CHANNEL-ABC flow not executed but remaining flow does get executed

CHANNEL-1 flow started
Aggregated
CHANNEL-1 flow ended

For File-4 //CHANNEL-ABC flow not executed also remaining flow doesn't get executed

CHANNEL-1 flow started
Aggregated

Sometime it skips calling DirectChannel and sometime it calls but return flow to some other channel. I want my flow to be executed in single thread, What is the best channel choice for above scenario if not DirectChannle? And Is above mentioned behavior is expected behavior, if Yes then How can I prevent it?

Update: I had added .log for explanation purpose only, Now I have updated the flow as well as observed result above.

1 Answers

You just faced the problem with the log() operator in the end of flow.

See more in docs: https://docs.spring.io/spring-integration/docs/current/reference/html/dsl.html#java-dsl-log:

When this operator is used at the end of a flow, it is a one-way handler and the flow ends.

That means your CHANNEL-ABC is going to have more than one subscriber with all those your configurations causing a default round-robin distribution strategy for DirectChannel: https://docs.spring.io/spring-integration/docs/current/reference/html/core.html#channel-implementations-directchannel

See this GH issue for more details: https://github.com/spring-projects/spring-integration/issues/3615

To fix the problem you just need to remove all those .log() operators in your flows for those CHANNEL-[N]. The one in the flow for CHANNEL-ABC is fine.

UPDATE

Recipient List Router approach:

    .routeToRecipients(r -> r
            .recipient("CHANNEL-ABC")
            .recipient("CHANNEL-3-POST-PROCESS"))
    .channel("CHANNEL-3-POST-PROCESS")
Related