Process a Small Import

Trace For Each, record per-item outcomes and trigger the same bounded work with a Scheduler.

Three records arrive together. The second has a zero quantity. We want an outcome for each record, including that rejection, without hiding what happened to the other two.

For Each runs the same processors over collection elements sequentially, which is what makes it a useful first import mechanism: one element can be traced all the way through the work before the next element begins. This example records outcomes in memory — it does not yet promise durable delivery.

Open this checkpoint in ACB

Stop the previous local application, then use File → Open Folder to open book/checkpoints/16-sequential-and-scheduled. Open src/main/mule/app.xml and use Flow List to select the named flow for each example. Start the controlled dependency in a separate terminal from the companion root with python3 book/stubs/server.py. Choose Run and Debug → Run Mule Application and wait for deployment. Save canvas edits and use Save and Hot-deploy to Local Runtime before repeating a request.

Run python3 book/run.py verify 16 from the companion root to exercise this running checkpoint with its synthetic fixtures. The verifier supplies requests and checks results; it does not start the ACB application. Keep the editor on this checkpoint while reading a failure so an old deployment cannot supply a misleading answer.

The bounded fixture is book/fixtures/import-records.json:

[
  {
    "orderId": "IMPORT-1",
    "qty": 4
  },
  {
    "orderId": "IMPORT-2",
    "qty": 0
  },
  {
    "orderId": "IMPORT-3",
    "qty": 2
  }
]

Example 041 — Process three records sequentially

The complete source is in checkpoints/16-sequential-and-scheduled/src/main/mule/app.xml.

Choose Flow List → process-records. The first Set Variable has Name outcomes and Value expression []. Expand For Each, then expand its nested Try.

Within Try, the Choice condition is payload.qty <= 0; its matching route raises APP:BAD_QUANTITY. The normal path’s next Set Variable updates outcomes using:

output application/java
---
vars.outcomes ++ [{orderId: payload.orderId, status: 'processed'}]

In Try → Error handler, On Error Continue matches APP:BAD_QUANTITY. Its Set Variable updates the same outcomes with:

output application/java
---
vars.outcomes ++ [{orderId: payload.orderId, status: 'rejected'}]

The final Transform Message sits after For Each. Its JSON payload script selects vars.outcomes. Inspect that nesting before moving a component: a handler around the whole loop would have a different continuation point.

The subflow initializes an empty outcomes array. The accumulator expressions explicitly use output application/java because they combine an in-memory array with JSON input; leaving the output representation implicit can produce a media-type inference error. Inside For Each, payload is the current record. ++ appends a one-element array to the accumulated outcomes. The Try handler owns only one iteration’s work: when quantity is invalid, it records rejected and permits the loop to continue.

After For Each, the scope restores the original collection as payload. Variable changes survive between iterations, which is why vars.outcomes contains all three classifications. The final transform selects that variable for the response. It would be a mistake to assume that setting each iteration’s payload automatically builds a new output collection.

POST the fixture to /lab/import. The result is IMPORT-1 processed, IMPORT-2 rejected and IMPORT-3 processed, in that order. Here processed means the local check accepted the quantity — not that an order reached a warehouse.

If the handler propagated instead, the failed iteration would fail the loop and prevent ordinary processing of subsequent elements. Choosing Continue therefore changes the import’s contract — the caller now needs the per-record result, since a successful request can contain a rejected record.

Input collection: IMPORT-1, 2, 3: Quantities 4, 0, 2. For Each + Try: Validate one record: Append its outcome. Outcome array: processed, rejected: processed. The Scheduler supplies a trigger; the subflow supplies the work.

Scheduler fixed frequency, start delay and local concurrency control

Scheduler fixed frequency, start delay and local concurrency control.

Let time start the same work

Nothing about that import needs a caller; it needs an occasion. A Scheduler is a source, like the HTTP Listener, but its trigger is time rather than an incoming request — which means the flow starts with nothing in hand, so give it an input explicitly. No HTTP body appears by itself.

Example 042 — Run the small import on a schedule

The complete source is in checkpoints/16-sequential-and-scheduled/src/main/mule/app.xml.

Choose Flow List → scheduled-import. Select the Scheduler source and inspect Fixed Frequency: Frequency 5, Start Delay 5, Time Unit SECONDS. Enable Disallow Concurrent Execution for this scheduler.

The next Set Payload uses expression readUrl("classpath://records.json", "application/json"). Flow Reference calls process-records. The final File Write uses Order_Files, Path scheduled-results.json, and Content expression output application/json --- payload.

Run the application in ACB and wait for a scheduled execution. Stop the runtime after the exercise so this five-second teaching job does not keep running.

The source waits five seconds, then runs at the configured frequency. It loads the same classpath fixture, calls the same subflow and serializes its result to scheduled-results.json. GET /lab/scheduled exposes that file for the local check. Wait for a newly written file after deploying; a file left by an earlier run is not evidence that this scheduler executed.

Disallow Concurrent Execution prevents overlapping executions of this scheduler in this runtime — it does not provide a distributed lease across independent replicas. Nor does a five-second period guarantee a job finishes in five seconds. If work takes longer, understand the scheduler’s missed-trigger behavior for the selected deployment. Scheduler reference.

This scheduler deliberately reprocesses a fixed fixture and overwrites a disposable result. A real synchronization job needs an eligibility rule, durable progress and replay-safe destination writes. Appendix D develops those responsibilities after the recovery and state chapters.

Size the loop to the request

For Each is reasonable for a short bounded collection. Its sequential nature can also protect a dependency from a burst of calls, although it is not a complete global rate limiter, because each request can still start its own loop.

Database pages should have an explicit ordering and bound. A fetch-size setting influences retrieval behavior — it does not replace a business limit or page cursor. If rows can change while pages are read, the query’s consistency and source-change contract determine whether records can be missed or repeated.

A bulk database operation takes parameter sets and reduces call overhead. Do not infer whole-import atomicity from the word bulk: inspect the driver result, transaction boundary and failure behavior. The local transaction in chapter 12 covered a deliberately small unit, not every page a scheduler might eventually process.

Try it

1. Explain the final payload. Remove the last transform in Process three records sequentially. What collection does For Each leave?

Show answer

The original input collection. The accumulated classifications are in vars.outcomes; they become the response only because the final transform selects them.

2. Stop on rejection. Change the per-record handler to Propagate. Which record cannot complete through the normal path?

Show answer

IMPORT-2 fails the loop, so IMPORT-3 does not reach its normal processing path in that execution. The source or enclosing flow must now own the failed import response.

3. Add a second runtime. Does the scheduler setting prevent both runtimes processing the same remote export?

Show answer

No. It prevents overlapping executions of this scheduler locally. Shared work requires a declared owner, supported distributed coordination or a destination that safely tolerates the overlap, plus a progress protocol.

Sequential execution made the state changes easy to follow. Parallel execution changes which work can overlap and how results return; those are useful choices only when the operations are independent.

Next: Run Independent Work in Parallel

Comments