Building a File-Queue Consumer With Idempotent Reprocessing
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
- 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.jsonextension into/root/q/inbox. There must be 20 files in inbox. - Work out the message structure and create
/root/q/schema.csv. The first line isfield,type,role. List every field in a normal message, and in therolecolumn, putidempotency-keyon the field that serves as the idempotency key andsequenceon the field used to judge ordering. - Create and run
/root/q/consume.py. It behaves as follows.- Move the files in
inboxone at a time toprocessing, then read them. - If JSON parsing fails or a required field (
msg_id,order_no,seq,amount) is missing, move it toerror. - If normal, load it into the sqlite DB
/root/q/ledger.db, in itsprocessedtable, and move it todone. - The
processedtable must havemsg_idas PRIMARY KEY or UNIQUE, and must have the columnsorder_no,seq,amount, andprocessed_at. - When the run ends,
processingmust be empty.
- Move the files in
- The run result must be as follows.
donehas 18,errorhas 2- The
processedtable row count is 15 after deduplication - Save the 3
msg_idvalues ignored as duplicates to/root/q/dup.txtin ascending order
- Organize the reasons for the messages that went to
errorin/root/q/error.csv. The first line isfile,reason.reasonisparseormissing-field. - Write
/root/q/order-check.sql. It is a query that finds, in theprocessedtable, cases within the sameorder_nowhereseqis 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 beOK. - Create
/root/q/lag.sh. It takes two arguments (큐디렉터리 임계치, that is, the queue directory and the threshold), prints one lineinbox=<n> processing=<n> error=<n>, and ends with a non-zero exit code if theinboxcount exceeds the threshold. - Create
/root/q/replay.sh. It takes one argument (a file name) and returns that file fromerrortoinbox. If the file is not inerror, do not change the queue state at all and end with a non-zero exit code.
Notes
- Checking the sqlite schema:
sqlite3 /root/q/ledger.db '.schema processed' - Ignoring duplicate inserts:
INSERT OR IGNOREorON CONFLICT DO NOTHING - Atomic move:
os.rename/mvwithin the same filesystem - Common mistake 1: sending parse-failure messages back to inbox, creating an infinite loop.
- Common mistake 2: preventing duplicates only with an application conditional. It is breached when run concurrently. The DB constraint is the last line of defense.
- Common mistake 3: not using
processingand reading directly frominbox. If you start two consumers, both process the same message.
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.
- Move the files in
inboxone at a time toprocessing, then read them. - If JSON parsing fails or a required field (
msg_id,order_no,seq,amount) is missing, move it toerror. - If normal, load it into the sqlite DB
/root/q/ledger.db, in itsprocessedtable, and move it todone. - The
processedtable must havemsg_idas PRIMARY KEY or UNIQUE, and must have the columnsorder_no,seq,amount, andprocessed_at. - When the run ends,
processingmust be empty.
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.
donehas 18,errorhas 2- The
processedtable row count is 15 after deduplication - Save the 3
msg_idvalues ignored as duplicates to/root/q/dup.txtin ascending order
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.