FDE PulseFDE jobs open 434New in the last 7 days 27
VI

The newspaper of the Forward Deployed Engineer

Guides

Adding AI to a client's existing event queue without rewriting the system

The client already has a queue that runs reliably. The FDE's job is to attach an AI consumer to it without charging a payment twice, letting failures pile up unnoticed or setting off an event storm.

Ảnh các kỹ sư đang trao đổi trước bảng trắng vẽ sơ đồ hệ thống trong văn phòng, hoặc màn hình laptop hiển thị code, gợi bối cảnh làm việc với hệ thống sẵn có của khách hàng.
Photo: Negative Space / CC0

In brief

  • A queue only guarantees that each message is delivered at least once, so the AI consumer must be idempotent from the first line of code.
  • Do not propose that the client move to event sourcing or CQRS. Attach the AI to the events they already publish.
  • Stopping event storms takes work on both sides: downstream services pass the correlation ID along, and the AI consumer skips tickets already labelled under the same correlation ID.
ShareLinkedInFacebookX

In your first week on site, you are handed a request that sounds simple enough: use an LLM to label every support ticket automatically as soon as it is created. The client’s engineering team opens the system diagram. Whenever a user submits a ticket, the UI pushes a message onto a queue, and several downstream services pick it up for processing.

This is the point at which many engineers start redrawing the whole architecture: add an event store, split out a read model, apply CQRS properly. That is usually the wrong move. The client did not hire you to rebuild their system. They need a new consumer that runs on the event stream they already have, runs safely and does not break what already works.

To do this well, you need to separate three concepts that are often used interchangeably, then write a consumer that avoids four familiar traps.

Three concepts, three levels of commitment

Event-driven is the simplest level. According to Microsoft’s Azure architecture documentation, a background job can be started when the UI or another job puts a message on a queue. The message carries data about an action that has just happened, such as a user placing an order. The ticketing system above is exactly this kind.

Event sourcing demands far more. Martin Fowler defines it as capturing every change to an application’s state as a sequence of events. Microsoft is more specific: the full sequence of actions is stored in an append-only store, and that store becomes the source of truth in place of a table of current state.

CQRS splits writes (commands) and reads (queries) into two separate data models so that each can be optimised independently. When the two sides use different stores, the common approach is for the write model to publish an event every time it updates the database, and for the read model to use those events to refresh its data.

One more confusion needs clearing up. Microsoft notes that an event store and a message broker should not be treated as the same thing. Kafka, or a similar queue, is only a distribution layer that delivers events to projections and external consumers. If a client says “we already have Kafka”, that does not mean they are doing event sourcing.

Why not to propose a re-architecture

Microsoft itself warns that event sourcing is expensive and rarely needed. For most systems, and for most components within a system, traditional data management is sufficient.

Fowler says much the same about CQRS: it is a significant mental leap for everyone involved, so it should only be undertaken when the benefit genuinely justifies it.

In discussing event sourcing, Microsoft cites one advantage: the code that publishes events is decoupled from the systems that subscribe to them. The same reasoning holds for an ordinary queue, even without event sourcing: the UI simply puts a message on the queue and does not need to know how many jobs are taking messages off it.

That means you can add an AI consumer without changing a single line in the application that publishes the events. For an FDE, this is a real advantage: the scope is small, there is little internal disruption, and you can roll back simply by switching the consumer off.

Example: a ticket-labelling consumer

Suppose each message is a TicketCreated with an event_id, ticket_id, correlation_id and the ticket’s content. The consumer below is a version you could actually deploy, not just demo:

def handle(msg):
    try:
        event = parse_ticket_created(msg.body)
    except SchemaError:
        msg.dead_letter(reason="schema")   # do not retry a broken message
        return

    if processed_events.exists(event.event_id):
        msg.complete()                     # duplicate message: skip
        return

    if labels.exists_for_correlation(event.ticket_id, event.correlation_id):
        msg.complete()                     # loop: ticket already labelled within the same business flow
        return

    label = classify_with_llm(event.text)  # safe to rerun, causes no harm

    with db.transaction():
        labels.upsert(event.ticket_id, label,
                      source_event=event.event_id,
                      correlation_id=event.correlation_id)
        processed_events.insert(event.event_id)
        outbox.add("TicketLabeled",
                   ticket_id=event.ticket_id,
                   correlation_id=event.correlation_id,
                   produced_by="ai-labeler")
    msg.complete()

Read the blocks from top to bottom and you will see that each one deals with a particular trap. The most important part is the three operations inside db.transaction(): saving the label, marking the event as processed and writing to the outbox table.

These three either all succeed or all fail. Microsoft recommends combining the Transactional Outbox pattern with idempotent consumers to keep the write and read models consistent, and the same principle applies unchanged to the AI consumer here.

Note that the LLM call sits outside the transaction. If the process dies after calling the model but before committing, the message will be redelivered and the model called once more. You pay a little extra in tokens, but no side effect is repeated. That is an acceptable trade-off.

Four traps and how the code avoids them

The first trap is duplicate delivery. Microsoft states plainly that a queue only guarantees at-least-once delivery, meaning the same message can arrive more than once.

If the handler is not idempotent, projections gradually drift away from the event stream, and side effects such as payments or notifications may run more than once. For an AI consumer with permission to email customers, that can become an incident within the first week.

The second trap is the poison message. A ticket may have empty content, or use an old schema that makes the parser fail. Microsoft points out that every message sitting in a dead-letter queue is a piece of unfinished work. If nobody watches that queue, failures pile up unnoticed, and the client only finds out when they see a lot of unlabelled tickets.

The third trap is the event storm. Choreography lets each service decide for itself when and how to act, usually through a broker using publish-subscribe. Microsoft warns that this design can inadvertently create feedback loops or event storms.

Imagine another service that listens for TicketLabeled and updates the ticket; that update publishes TicketCreated again, and the AI consumer runs once more. The new event is published by the other service, so it will not carry produced_by="ai-labeler" on its own, and checking event_id will not stop it either, because it is a new event.

The fix is to agree with the client’s team that every downstream service passes the correlation_id on to the events it publishes.

Then the labels.exists_for_correlation step in the code earns its place: a ticket already labelled under the same correlation_id is skipped, and the loop stops on the second pass.

Without that forwarding on the client’s side, the check has nothing to compare against.

The fourth trap is stale data. When the read and write stores are separate, the system is only eventually consistent, so reads may not yet reflect the latest change.

If your agent reads customer history from a read model to decide on a refund, it may decide on data that is a few seconds old. If the decision needs the freshest data, take the information directly from the event payload, or ask the client’s team how far the read model typically lags.

Correlation IDs: so you can still answer “what happened?”

Without a central orchestrator, Microsoft notes, no single component can see an entire business process as it unfolds. That is why correlation IDs and distributed tracing are needed. When a customer complains that their ticket was mislabelled, you need to trace back from that label to the original event, then to the model call and the prompt that was used.

The method is simple. Every event you publish must carry the correlation_id of the event that triggered it. Every log line from the consumer should record this ID too. If the client has no such ID yet, propose adding one before the AI goes into production. It is a small change, but it pays off in the first incident investigation.

The order of work on site

First, establish which of the three levels the client is at. Microsoft suggests event sourcing fits best when an application already uses events as a natural part of how it operates.

Many clients will have only a queue, and that is enough. Next, ask for a sample of real messages, read the schema and ask whether they have ever changed it.

Then run the consumer in shadow mode: it still labels tickets but writes only to a separate table and publishes no events. Only when the number of messages in the dead-letter queue is stable and the labels meet the bar do you switch on the outbox. This lets you separate two concerns: the quality of the model and the safety of the system.

If you are preparing to move into an FDE role, put exactly these details on your CV. “Integrated an LLM into a Kafka pipeline” says very little.

“Wrote an idempotent consumer using an outbox, with dead-letter alerting and correlation IDs, run in shadow mode for two weeks before publishing events” tells a hiring manager you understand the risks of working on someone else’s system.

When a job description asks for integration with existing systems or experience with message queues, be ready to walk through exactly one example like this in the interview, including how you handled duplicate delivery.

What impresses in a demo is the model labelling accurately. What decides whether the client keeps trusting you is that the consumer stays safe when the same message arrives a second time.

6 sources
Read next on the roadmap · Stage 2: Broad engineeringLatency, throughput, scalability: answering the client's infrastructure team with measurementsWhen the client's infrastructure team asks "what's the p99, and how many requests per second can it take?", answering "it runs pretty fast" will cost you their trust in the very first meeting.