TT Lab
Get started
Learn Learning paths Courses

System Integration (EAI)

Building a File-Queue Consumer With Idempotent Reprocessing

Continue in TT Lab

Goal

You build a consumer on a directory-based queue, isolate parse failures, guarantee idempotency with a DB constraint, and implement lag monitoring and reprocessing.

Why it matters

In asynchronous integration, duplicates are the default, not the exception. "Exactly-once" is not the magic of the delivery layer but the result of "at-least-once delivery + receiver-side deduplication." So the consumer must always be written assuming duplicates. And if you put a parse-failure message in the retry queue, that message keeps failing at the head and blocks the normal messages behind it (a poison message). A blocked queue is a business stop. Implementing these two by hand builds your feel for asynchronous design.

Steps

  1. Create /root/q/inbox, /root/q/processing, /root/q/done, and /root/q/error, and copy all the files of /opt/lab/fixtures/eai/queue/inbox/ that have a .json extension into /root/q/inbox. There must be 20 files in inbox.
  2. Work out the message structure and create /root/q/schema.csv. The first line is field,type,role. List every field in a normal message, and in the role column, put idempotency-key on the field that serves as the idempotency key and sequence on the field used to judge ordering.
  3. Create and run /root/q/consume.py. It behaves as follows.
    • Move the files in inbox one at a time to processing, then read them.
    • If JSON parsing fails or a required field (msg_id, order_no, seq, amount) is missing, move it to error.
    • If normal, load it into the sqlite DB /root/q/ledger.db, in its processed table, and move it to done.
    • The processed table must have msg_id as PRIMARY KEY or UNIQUE, and must have the columns order_no, seq, amount, and processed_at.
    • When the run ends, processing must be empty.
  4. The run result must be as follows.
    • done has 18, error has 2
    • The processed table row count is 15 after deduplication
    • Save the 3 msg_id values ignored as duplicates to /root/q/dup.txt in ascending order
  5. Organize the reasons for the messages that went to error in /root/q/error.csv. The first line is file,reason. reason is parse or missing-field.
  6. Write /root/q/order-check.sql. It is a query that finds, in the processed table, cases within the same order_no where seq is duplicated or missing. Save the result (0 rows if there is no problem) to /root/q/order-result.txt. If there is no problem, the first line of the file must be OK.
  7. Create /root/q/lag.sh. It takes two arguments (큐디렉터리 임계치, that is, the queue directory and the threshold), prints one line inbox=<n> processing=<n> error=<n>, and ends with a non-zero exit code if the inbox count exceeds the threshold.
  8. Create /root/q/replay.sh. It takes one argument (a file name) and returns that file from error to inbox. If the file is not in error, do not change the queue state at all and end with a non-zero exit code.

Notes

Set up the queue directories

Create /root/q/inbox, /root/q/processing, /root/q/done, and /root/q/error, and copy all the files of /opt/lab/fixtures/eai/queue/inbox/ that have a .json extension into /root/q/inbox. There must be 20 files in inbox.

Create the four sections inbox/processing/done/error. Remember that they must be on the same filesystem for moves to be atomic.

Work out the message structure

Work out the message structure and create /root/q/schema.csv. The first line is field,type,role. List every field in a normal message, and in the role column, put idempotency-key on the field that serves as the idempotency key and sequence on the field used to judge ordering.

Open a message sample and organize the fields. Pay attention to which field could become the idempotency key.

Implement and run the consumer

Create and run /root/q/consume.py. It behaves as follows.

A parse failure fails the same way however often you retry. Send such messages to error immediately so the normal messages behind them are not blocked. Record the processing history in the same transaction as the business processing.

Verify idempotency

The run result must be as follows.

Blocking duplicates with a DB constraint is safer than filtering them in code. Even if the application has a bug, the constraint cannot be breached.

Classify failure reasons

Organize the reasons for the messages that went to error in /root/q/error.csv. The first line is file,reason. reason is parse or missing-field.

Try dividing the reasons into "format errors" and "business errors." The former need a fix at the source, and the latter can be re-injected after supplementing the data.

Verify ordering guarantees

Write /root/q/order-check.sql. It is a query that finds, in the processed table, cases within the same order_no where seq is duplicated or missing. Save the result (0 rows if there is no problem) to /root/q/order-result.txt. If there is no problem, the first line of the file must be OK.

Check with SQL whether the messages with the same key were processed in sequence order. You must compare the stored sequence numbers, not the processing times.

Lag monitoring script

Create /root/q/lag.sh. It takes two arguments (큐디렉터리 임계치, that is, the queue directory and the threshold), prints one line inbox=<n> processing=<n> error=<n>, and ends with a non-zero exit code if the inbox count exceeds the threshold.

Taking the queue directory as an argument makes it reusable. If you signal with an exit code when the threshold is exceeded, you can attach it to cron or a monitoring tool.

Reprocessing script

Create /root/q/replay.sh. It takes one argument (a file name) and returns that file from error to inbox. If the file is not in error, do not change the queue state at all and end with a non-zero exit code.

Reprocessing is not putting just anything back. If it receives a message ID that does not exist, it must do nothing and fail.