Fluxo de bloco de integração Spring com Loop e vários HTTP e SQS assíncronos

Aug 26 2020

Eu tenho um fluxo que

1. Starts with a config map -> MainGateway.start(configMap) -> void
2. Splits map into multiple messages per entry
3. For every config entry do the following using an orchestrator java class:
   BEGIN LOOP (offset and limit)
      Data d = HTTPGateway.getData();
      PublishGateway.sendMessage(d); -> Send to 2 SQS queues    
   END LOOP

Requisito Tenho que agendar este fluxo via cron. Uma opção é fornecer um endpoint HTTP que iniciará o fluxo. Mas então a segunda solicitação HTTP deve esperar / tempo limite / erro até que a primeira seja concluída.

Questão Eu estava procurando uma barreira para implementar o bloqueio do thread de fluxo até que ele seja concluído e tenha apenas um processador http de thread único, portanto, ao mesmo tempo, apenas 1 solicitação é processada e posso saber quando o fluxo está concluído. (O LOOP termina para todos os objetos de entrada de configuração e todas as mensagens para SQS são confirmadas). Como posso conseguir isso? Eu tenho um loop e estou usando o canal pub-sub com executores para configurações paralelas e despacho SQS paralelo.

Eu reduzi o XML configabaixo para maior clareza.

   <!-- Bring in list of Configs to process -->
    <int:gateway service-interface="Gateway"
                 default-request-channel="configListChannel" />

    <int:chain input-channel="configListChannel" output-channel="configChannel">
        <!-- Split the list to one instance of config per message -->
        <int:splitter/>
        <int:filter expression="payload.enablePolling" />
    </int:chain>

    <!-- Manually orchestrate a loop to query a system as per config and publish messages to SQS -->
    <bean class="Orchestrator" id="orchestrator" />
    <int:service-activator ref="orchestrator" method="getData" input-channel="configChannel" />

    <!-- The flow from this point onwards is triggered inside a loop controlled by the Orchestrator
         The following Gateway calls are inside Orchestrators loop -->

    <!-- Create a Http request from the Orchestrator using a Gateway -->
    <int:gateway service-interface="HttpGateway">
        <int:method name="getData"
                    request-channel="requestChannel"
                    payload-expression="#args[0]">
        </int:method>
    </int:gateway>

    <!-- Transform request object to json and invoke Http endpoint -->
    <int:chain input-channel="requestChannel" id="httpRequestChain">
        <int:object-to-json-transformer />
        <int-http:outbound-gateway url-expression="headers['config'].url"
                                   http-method="POST"
                                   expected-response-type="java.lang.String"
        />
    </int:chain>

    <!-- Publish Messages to Outbound Gateway -->
    <task:executor id="executor" pool-size="5" />
    <int:publish-subscribe-channel id="publishChannel" task-executor="executor" />
    <int:gateway service-interface="PublishGateway" >
        <int:method name="publishToOutbound" payload-expression="#args[0]" request-channel="publishChannel" />
    </int:gateway>


    <!-- Route to System A SQS with transformations (omitted here)-->
    <int-aws:sqs-outbound-channel-adapter sqs="amazonSQS" channel="publishChannel" queue="system-a-queue" success-channel="successChannel" failure-channel="errorChannel"/>

    <!-- Route to System B SQS with transformations (omitted here)-->
    <int-aws:sqs-outbound-channel-adapter sqs="amazonSQS" channel="publishChannel" queue="system-b-queue" success-channel="successChannel" failure-channel="errorChannel"/>

    <int:logging-channel-adapter logger-name="sqsCallbackLogger" log-full-message="true" channel="successChannel" />

Nesse ínterim, estou tentando adaptar o exemplo da barreira ABC de spring-integration-samplespara o meu caso de uso.

Respostas

1 ArtemBilan Aug 26 2020 at 17:54

Como você apontou em seu comentário, uma aggregatorabordagem pode ser usada em sua solução.

Dessa forma, você agregou os resultados dessas solicitações SQS paralelas e aguardou uma resposta de agregação no solicitante original. Dessa forma, ele realmente será bloqueado, mesmo se os componentes internos do seu fluxo ainda estiverem simultâneos. Você chama um gateway e a resposta para ele será do agregador.