Transaction Log Tailing is the way to get the events in your Transactional Outbox out to a message broker without writing a polling loop.
Intro
Transaction Log Tailing is the way to get the events in your Transactional Outbox out to a message broker without writing a polling loop. Every relational database already records each committed change in its transaction log (the write-ahead log, WAL, in PostgreSQL; the transaction log and CDC change tables in SQL Server). A log-tailing process, usually Debezium, reads that log, notices the new outbox rows, and publishes them to Kafka or Azure Event Hubs. Your application only does one thing: insert the business row and the outbox row in one local transaction. The database log becomes the reliable, ordered source of truth for publishing, with no query load on your tables and very low latency.
Why we need this
In Day 12 we solved the dual-write problem: ClaimService saves the claim and an OutboxMessages row in the same ACID transaction, so we can never save a claim without also having recorded the event. But an outbox row sitting in a table is not yet an event on a broker. Someone has to move it there, and that “someone” must be:
- Reliable: no committed event may be lost, even if the relay crashes mid-way.
- Ordered: events for the same claim must reach the broker in commit order (
ClaimSubmittedbeforeClaimApproved). - Fast: the fraud-check service should see
ClaimSubmittedin milliseconds to seconds, not after the next 30-second poll. - Cheap for the database: publishing must not steal CPU and IOPS from the claims workload.
The database log already gives you all four properties for free, because it is what the database itself uses for crash recovery and replication. Tailing it means we reuse a mechanism that is battle-tested and already in commit order.
What problem it solves
Problem from the topic list: polling creates query overhead and propagation latency.
In a system that uses a polling outbox publisher (Day 14), a background job runs SELECT ... FROM OutboxMessages WHERE ProcessedAt IS NULL ORDER BY Id every N seconds. That creates concrete pains:
- Latency floor: an event is never delivered faster than the poll interval. Poll every 5 seconds and your average delay is about 2.5 seconds.
- Constant query load: the query runs 17,280 times a day per instance at a 5-second interval even when there is nothing to send. Lower the interval to fix latency and the load rises.
- Index and lock contention: the “unprocessed” query and the
UPDATE ... SET ProcessedAtcompete with inserts on the same hot table. Under heavy claim intake, this shows up as lock waits. - Ordering hazards with multiple pollers: two relay instances polling the same table can publish out of order unless you add locking or partitioning.
- Table bloat: marking rows processed keeps them around; deleting them adds more write load.
With log tailing, the relay is a passive reader of a sequential file. Nothing is queried, nothing is marked, and the outbox table can be write-only (insert, then periodic cleanup).
When it is needed (and when it is NOT)
Good fit:
- You already use the Transactional Outbox and event volume or latency requirements make polling painful (thousands of events per minute, or sub-second delivery needs).
- You run PostgreSQL or SQL Server (or MySQL, Oracle, MongoDB) with a supported CDC connector and you have admin rights to enable it.
- You want the same pipeline to feed several consumers (fraud, notifications, analytics) through Kafka or Event Hubs.
- You already operate Kafka Connect or can run a container for Debezium Server.
Not a good fit:
- Low volume (a few hundred events a day). A 5-second polling publisher in your own .NET worker is simpler to run, debug, and explain. Log tailing adds a new moving part (Debezium, connector offsets, replication slots).
- The database does not expose a log to you: managed platforms that block logical replication or CDC, or Azure SQL Database on a tier that does not support CDC (see section 9). Use the Polling Publisher (Day 14).
- Team cannot operate it. A stuck replication slot can fill the disk of your production database. If nobody owns monitoring of it, do not adopt it yet.
- You need to capture changes without modifying the app (legacy monolith). That is a valid CDC use case, but it is Change Data Capture of table rows, not the outbox pattern. It leaks your table schema to other services as a contract, which is usually a coupling problem.
How to identify the problem (key signals)
You probably need log tailing (over polling) when you see several of these:
- Event delivery latency tracks the poll interval. Traces show a consistent gap equal to the polling period between
ClaimSubmittedcommitted andClaimSubmittedconsumed. - Database CPU or query stats show the outbox query near the top of
pg_stat_statementsor SQL Server Query Store by execution count, even though it returns zero rows most of the time. - Lock waits or deadlocks on the outbox table during peak intake (for example after a storm or catastrophe event that triggers many claims).
- Outbox table is growing because cleanup jobs cannot keep up, and the
Unprocessedindex gets bigger and slower. - Out-of-order events on the consumer side after you scaled the relay to two instances.
- Team keeps lowering the poll interval from 30 seconds to 5 seconds to 1 second. That is a smell that polling has reached its limit.
- Business complaint: “The adjuster dashboard shows the claim as Submitted for a minute after the customer got the approval email.”
Flow Diagram
The relay tails the database log instead of polling the outbox table.
flowchart LR APP["ClaimSubmitter"] -- "one transaction" --> DB[("PostgreSQL: claims + outbox")] DB --> WAL["WAL / transaction log"] WAL -- "replication slot" --> DBZ["Debezium + Outbox Event Router"] DBZ -- "key = aggregate_id" --> EH{{"Event Hubs / Kafka topic"}} EH --> FR["Fraud service"] EH --> NO["Notification service"] EH --> RM["Read-model service"]Level 1: Beginner
Analogy. Think of a bank branch. Every deposit and withdrawal is written in the branch’s ledger (the transaction log). Polling is like a clerk who every five minutes walks to the ledger, scans the whole book for new entries, and phones head office. Log tailing is like a second clerk who sits at the ledger and reads each new line as the first clerk writes it, immediately telephoning head office. The ledger is never re-scanned, and no entry is missed because the reader remembers exactly which line it read last (its offset).
The three moving parts:
- Your service writes the business row and an outbox row in one transaction.
- A log reader (Debezium) reads the database log and turns each new outbox row into a message.
- A broker (Kafka / Event Hubs) stores the message for consumers.
Minimal working example. Outbox table (PostgreSQL) and a transaction that writes both rows. The table is intentionally shaped for Debezium’s outbox event router.
-- PostgreSQLCREATE TABLE claims ( id uuid PRIMARY KEY, policy_no text NOT NULL, status text NOT NULL, amount numeric(12,2) NOT NULL, created_at timestamptz NOT NULL DEFAULT now());
CREATE TABLE outbox ( id uuid PRIMARY KEY, aggregate_type text NOT NULL, -- e.g. 'claim' -> used for topic routing aggregate_id text NOT NULL, -- e.g. claim id -> used as Kafka message key type text NOT NULL, -- e.g. 'ClaimSubmitted' payload jsonb NOT NULL, created_at timestamptz NOT NULL DEFAULT now());// .NET 10, C# 14 - plain Npgsql, showing that the app only does ONE local transactionusing Npgsql;using System.Text.Json;
public sealed class ClaimSubmitter(NpgsqlDataSource db){ public async Task<Guid> SubmitAsync(string policyNo, decimal amount, CancellationToken ct) { var claimId = Guid.NewGuid(); await using var conn = await db.OpenConnectionAsync(ct); await using var tx = await conn.BeginTransactionAsync(ct);
await using (var cmd = new NpgsqlCommand( "INSERT INTO claims(id, policy_no, status, amount) VALUES ($1, $2, 'Submitted', $3)", conn, tx)) { cmd.Parameters.AddWithValue(claimId); cmd.Parameters.AddWithValue(policyNo); cmd.Parameters.AddWithValue(amount); await cmd.ExecuteNonQueryAsync(ct); }
await using (var cmd = new NpgsqlCommand( "INSERT INTO outbox(id, aggregate_type, aggregate_id, type, payload) " + "VALUES ($1, 'claim', $2, 'ClaimSubmitted', $3)", conn, tx)) { cmd.Parameters.AddWithValue(Guid.NewGuid()); cmd.Parameters.AddWithValue(claimId.ToString()); cmd.Parameters.Add(new NpgsqlParameter { Value = JsonSerializer.Serialize(new { claimId, policyNo, amount }), NpgsqlDbType = NpgsqlTypes.NpgsqlDbType.Jsonb }); await cmd.ExecuteNonQueryAsync(ct); }
await tx.CommitAsync(ct); // both rows are in the WAL only if this commits return claimId; }}Notice there is no publishing code in the service. Debezium takes it from here.
Level 2: Intermediate
6.1 The Debezium connector (PostgreSQL)
Debezium runs as a Kafka Connect connector (or as Debezium Server, a standalone process that can write straight to Azure Event Hubs). The connector registers a logical replication slot, reads decoded changes through the built-in pgoutput plugin, and applies the Outbox Event Router transformation so each outbox row becomes a clean business event rather than a raw CDC envelope.
{ "name": "claims-outbox-connector", "config": { "connector.class": "io.debezium.connector.postgresql.PostgresConnector", "database.hostname": "claims-pg.postgres.database.azure.com", "database.port": "5432", "database.user": "debezium", "database.password": "${file:/secrets/db.properties:password}", "database.dbname": "claims", "database.sslmode": "require", "topic.prefix": "claims-db", "plugin.name": "pgoutput", "slot.name": "claims_outbox_slot", "publication.name": "claims_outbox_pub", "publication.autocreate.mode": "filtered", "table.include.list": "public.outbox", "tombstones.on.delete": "false",
"transforms": "outbox", "transforms.outbox.type": "io.debezium.transforms.outbox.EventRouter", "transforms.outbox.table.field.event.id": "id", "transforms.outbox.table.field.event.key": "aggregate_id", "transforms.outbox.table.field.event.type": "type", "transforms.outbox.table.field.event.payload": "payload", "transforms.outbox.route.by.field": "aggregate_type", "transforms.outbox.route.topic.replacement": "claims.${routedByValue}.events", "transforms.outbox.table.fields.additional.placement": "type:header:eventType" }}What this gives you:
- Rows in
outboxwithaggregate_type = 'claim'go to the topicclaims.claim.events. - The Kafka message key is
aggregate_id, so all events for one claim land on one partition and stay in order. - The message value is the
payloadJSON, andidis delivered as a header so consumers can de-duplicate. - The router only reacts to inserts; delete events (from cleanup) are ignored, so cleaning up the table does not produce events.
6.2 Cleaning up the outbox
Since nothing marks rows as processed, cleanup is a separate concern. In PostgreSQL, a common technique from the Debezium documentation is to insert and delete the row in the same transaction: the log still contains the insert, so Debezium sees it, but the table never grows. A simpler and safer alternative for beginners is a nightly job:
DELETE FROM outbox WHERE created_at < now() - interval '2 days';Keep the retention longer than your worst realistic connector outage so you can investigate, but remember the log, not the table, is what guarantees delivery.
6.3 The consumer: still at-least-once
Log tailing gives at-least-once delivery. After a connector restart, Debezium can re-deliver events since the last committed offset. The consumer in FraudService must therefore be idempotent (Day 17). Minimal .NET consumer using Confluent.Kafka:
// dotnet add package Confluent.Kafkausing Confluent.Kafka;using System.Text;
public sealed class ClaimEventsWorker(IServiceScopeFactory scopes, ILogger<ClaimEventsWorker> log) : BackgroundService{ protected override Task ExecuteAsync(CancellationToken stoppingToken) => Task.Run(() => Run(stoppingToken), stoppingToken);
private async Task Run(CancellationToken ct) { var cfg = new ConsumerConfig { BootstrapServers = "kafka:9092", GroupId = "fraud-service", EnableAutoCommit = false, // commit only after we processed successfully AutoOffsetReset = AutoOffsetReset.Earliest };
using var consumer = new ConsumerBuilder<string, string>(cfg).Build(); consumer.Subscribe("claims.claim.events");
while (!ct.IsCancellationRequested) { var result = consumer.Consume(ct); var eventId = Encoding.UTF8.GetString(result.Message.Headers.GetLastBytes("id"));
using var scope = scopes.CreateScope(); var handler = scope.ServiceProvider.GetRequiredService<IClaimEventHandler>();
// Idempotency: the handler inserts eventId into a ProcessedEvents table // with a unique key inside the same DB transaction as its own work. await handler.HandleAsync(eventId, result.Message.Key, result.Message.Value, ct);
consumer.Commit(result); log.LogInformation("Handled {EventId} for claim {ClaimId}", eventId, result.Message.Key); } }}6.4 Where Angular fits
The Angular adjuster dashboard never talks to the broker. It reads a read model that a consumer keeps updated. Because delivery is asynchronous, the UI must accept eventual consistency: after submitting a claim, show “Submitted, being reviewed” and refresh the status.
// Angular (standalone component, signals) - polls the read-model API every 3 secondsimport { Component, inject, signal } from '@angular/core';import { HttpClient } from '@angular/common/http';import { interval, switchMap, startWith, map } from 'rxjs';import { toSignal } from '@angular/core/rxjs-interop';
@Component({ selector: 'app-claim-status', template: `<p>Claim status: {{ status() ?? 'Loading...' }}</p>`})export class ClaimStatusComponent { private http = inject(HttpClient); claimId = '...'; // supplied by route or input in real code
status = toSignal( interval(3000).pipe( startWith(0), switchMap(() => this.http.get<{ status: string }>(`/api/claims/${this.claimId}/status`)), map(r => r.status), ), { initialValue: undefined as string | undefined } );}The point is the pattern: poll or push the read model, and never assume the event has already been processed.
6.5 SQL Server variant
SQL Server does not expose a logical decoding API. Debezium’s SQL Server connector relies on SQL Server CDC: the engine’s capture job reads the transaction log and writes changes into cdc.<capture_instance>_CT tables, and Debezium reads those change tables. Enabling it:
-- SQL Server / Azure SQL (needs sysadmin or db_owner depending on platform)EXEC sys.sp_cdc_enable_db;
EXEC sys.sp_cdc_enable_table @source_schema = N'dbo', @source_name = N'OutboxMessages', @role_name = NULL; -- no gating role; restrict via permissions in real systemsIt is still “log tailing” in effect, but be honest with your team: for SQL Server there is one extra hop (log to CDC tables to Debezium), and CDC’s capture job and cleanup job add some load and have their own retention window (default 3 days).
Level 3: Advanced
7.1 Performance and scalability
- Throughput: one connector task reads one log stream sequentially. A single PostgreSQL connector can handle thousands of events per second on modest hardware; the practical limit is usually the broker or downstream consumers. Verify with your own load test rather than trusting a number.
- Partitioning: use
aggregate_idas the message key. Never key byaggregate_typealone, or one hot partition receives everything. - Batching: tune
max.batch.sizeandmax.queue.sizeon the connector when snapshots or bursts cause memory pressure. - Keep the outbox narrow: keep the
payloadreasonably small (the log records the full row). Store large documents in Blob Storage and put a reference in the event (claim-check pattern). - Replica: PostgreSQL logical decoding runs on the primary (a standby cannot host the slot in most managed configurations). Plan for the extra WAL sender load.
7.2 Failure modes (the ones that page you at 3 a.m.)
| Failure mode | What happens | Prevention / response |
|---|---|---|
| Connector down for a long time (PostgreSQL) | The replication slot keeps WAL segments the connector has not confirmed. WAL piles up and can fill the disk, taking the whole database down. | Alert on pg_replication_slots (active, confirmed_flush_lsn) and retained WAL size. Set max_slot_wal_keep_size (PostgreSQL 13+) so the slot is invalidated before the disk fills, and accept a re-snapshot. Drop unused slots. |
| Connector down longer than CDC retention (SQL Server) | Change rows are cleaned up before Debezium reads them, so events are lost from the stream. | Alert on connector lag. Set CDC retention above the worst realistic outage. Be ready to re-snapshot and republish. |
| Duplicate delivery after restart | The consumer sees the same event again. | Idempotent consumer keyed on the event id header (Day 17). |
| Schema change on the outbox table | Connector can fail or emit changed shapes. | Treat the outbox schema as a contract. Add nullable columns only. Test migrations against the connector in staging. |
| Failover of the database | Slot or offsets may not carry over. | On PostgreSQL, use a platform version that supports slot synchronization for HA and test the failover. Otherwise expect a resnapshot or gap analysis after failover. |
| Poison payload (invalid JSON) | Connector or consumer stops on one bad record. | Validate payloads at write time. Use a dead-letter topic on the consumer side. |
| Broker unavailable | Connector retries and stops advancing offsets; the source log keeps growing. | Same WAL/retention alerts as above. The outbox pattern means no business data is lost, only delayed. |
7.3 Security
- Create a dedicated
debeziumdatabase user with the least privileges:REPLICATIONplusSELECTon the outbox table (PostgreSQL), ordb_datareaderon the CDC schema (SQL Server). - Use TLS from connector to database (
sslmode=requireor stricter) and to the broker. - Keep credentials out of connector JSON: use Kafka Connect config providers or Key Vault-backed environment variables.
- Do not put PII in the payload without thinking. The event goes to a broker that many teams may read. Put a
claimId, not the claimant’s full medical narrative. Encrypt or omit sensitive fields. - Restrict who can read the CDC change tables and the slot. They contain the raw data of the tracked tables.
7.4 Common mistakes
- Capturing business tables instead of the outbox. This publishes your internal schema and creates hidden coupling. Capture only the outbox.
- Forgetting slot monitoring (PostgreSQL) and discovering the problem when the disk is full.
- Assuming exactly-once. It is at-least-once; the consumer must be idempotent.
- Wrong message key, breaking per-claim ordering.
- Running two connectors with the same slot name or restarting with a new slot name and silently skipping or replaying history.
- Not testing the snapshot behavior. By default Debezium takes an initial snapshot of included tables. For an outbox you usually want
snapshot.mode=no_data(orneverin older versions) so old, already-published rows are not published again. Check the option names for your Debezium version. - Using the outbox as an event store. Rows are deleted; it is a transient buffer (compare with Event Sourcing, Day 9).
Level 4: Expert and Architect view
8.1 Options compared
| Criterion | Transaction Log Tailing (Debezium/CDC) | Polling Publisher (Day 14) | Dual write (anti-pattern) | Managed change feed (e.g. Cosmos DB change feed) |
|---|---|---|---|---|
| Delivery guarantee | At-least-once, ordered per key | At-least-once, ordering needs care | None (partial failures) | At-least-once, ordered per partition key |
| Latency | Sub-second to a few seconds | Poll interval | Lowest but unsafe | Low |
| Load on source DB | Very low (sequential log read) | Continuous queries and updates | None extra | Handled by the platform |
| Operational complexity | High (connector, slots, offsets, broker) | Low (a .NET worker) | Lowest | Medium |
| Database support | Needs CDC/logical decoding support | Any database | Any | Only that database |
| Risk of taking down the DB | Yes if slot lag ignored | Low | Data inconsistency instead | Low |
| Best for | High volume, low latency, several consumers | Modest volume, restricted platforms | Never in production | Cosmos DB based services |
8.2 Combines with
- Transactional Outbox (Day 12): log tailing is one of two ways to relay the outbox; it is not a standalone pattern.
- Idempotent Consumer (Day 17): required, because relay delivery is at-least-once.
- Saga (Day 7): sagas rely on reliable event or command delivery between steps; log-tailed outbox events are a good backbone.
- Domain Event (Day 11): the domain event becomes the outbox row’s payload.
- CQRS (Day 8) and API Composition (Day 10): read models are refreshed by consumers of these events.
- Messaging (Day 16): the broker at the end of the pipeline.
- Health checks and metrics (Days 33, 36): connector lag and slot size are first-class operational metrics.
8.3 ADR
ADR-013: Use log-based relay (Debezium) for publishing Claims outbox events
- Status: Proposed
- Context: The Claims service publishes
ClaimSubmitted,ClaimApproved, andClaimRejectedevents using the Transactional Outbox. The current polling publisher runs every 5 seconds per instance and contributes measurable load on the claims database during intake peaks. Fraud, Notification, and Reporting services need events within 2 seconds and in per-claim order. - Decision: Capture only the
outboxtable from the Claims PostgreSQL database with Debezium (logical decoding,pgoutput), route with the Outbox Event Router to topicsclaims.<aggregate_type>.eventskeyed byaggregate_id, and publish to Azure Event Hubs through its Kafka endpoint. Consumers are idempotent by event id. - Consequences:
- (+) Latency drops from the poll interval to about a second; query load on the outbox table disappears.
- (+) Consistent, ordered delivery for several downstream services.
- (-) New components to operate: Debezium runtime, replication slot, connector offsets.
- (-) Risk of WAL growth if the connector is down; requires alerting and
max_slot_wal_keep_size. - (-) Team needs runbooks for connector restart, re-snapshot, and failover.
- Alternatives considered: keep polling publisher (rejected: latency and load); dual write (rejected: inconsistency); Azure SQL/SQL Server CDC (not chosen because the claims store is PostgreSQL).
- Review trigger: revisit if event volume stays below a few hundred per hour, since polling would then be cheaper to run.
Azure implementation
9.1 Services involved
| Role | Azure service | Notes |
|---|---|---|
| Source database (PostgreSQL) | Azure Database for PostgreSQL - Flexible Server | Set wal_level = logical (requires a restart), and raise max_replication_slots and max_wal_senders if needed. |
| Source database (SQL) | Azure SQL Database or SQL Managed Instance | Supports CDC. On Azure SQL Database with the DTU model, CDC needs the S3 tier or higher (Basic, S0, S1, S2 are not supported); vCore tiers are supported. Confirm current limits in the Microsoft Learn CDC article before you choose a tier. |
| Broker | Azure Event Hubs (Kafka endpoint) | The Kafka protocol endpoint requires the Standard tier or above (Basic does not support Kafka). Premium and Dedicated add isolation and higher throughput. Alternatively use Confluent Cloud or a self-hosted Kafka cluster. |
| Relay runtime | Azure Container Apps or AKS | Run Kafka Connect with the Debezium connector, or Debezium Server (which has an Azure Event Hubs sink). One replica only for a given slot. |
| Secrets | Azure Key Vault | Database and broker credentials, referenced through managed identity where the runtime supports it. |
| Monitoring | Azure Monitor, Log Analytics, Application Insights | Connector logs, database metrics, Event Hubs metrics, alerts. |
| Consumers | Azure Container Apps / AKS / Azure Functions (Event Hubs trigger) | Idempotent handlers. |
9.2 How to configure
- PostgreSQL Flexible Server: in Server parameters set
wal_leveltological, reviewmax_replication_slotsandmax_wal_senders, restart the server. Create thedebeziumrole and grant it the replication attribute per the Azure documentation (on the managed service, use the platform-provided replication role mechanism rather than superuser). Create theoutboxtable and, if you prefer, the publication in advance. - Event Hubs: create a namespace on Standard or higher. In Event Hubs, a Kafka topic is an event hub, so pre-create
claims.claim.eventsand any Debezium internal topics (offsets, config, status for Kafka Connect, and schema history if used) rather than relying on auto-creation; several Debezium and Kafka Connect setups on Event Hubs fail when they try to create topics through the Kafka admin API. Kafka clients connect to<namespace>.servicebus.windows.net:9093withSASL_SSL, mechanismPLAIN, username$ConnectionString, and the connection string (or a Microsoft Entra ID / OAuth setup) as the password. - Relay: deploy Kafka Connect (or Debezium Server) to Azure Container Apps with min and max replicas both set to 1. Load the connector JSON from section 6.1. Put secrets in Key Vault and inject them as secret references.
- Networking: use private endpoints for the database and Event Hubs, or VNet integration and firewall rules, so the relay reaches both privately.
- Consumers: use consumer groups per service (
fraud-service,notification-service).
Debezium Server sink example (properties, simplified):
debezium.sink.type=eventhubsdebezium.sink.eventhubs.connectionstring=${EVENTHUBS_CONNECTION_STRING}debezium.sink.eventhubs.hubname=claims.claim.eventsdebezium.source.connector.class=io.debezium.connector.postgresql.PostgresConnectordebezium.source.plugin.name=pgoutputdebezium.source.table.include.list=public.outboxdebezium.source.topic.prefix=claims-dbdebezium.source.offset.storage.file.filename=/data/offsets.datDebezium Server keeps offsets in the storage you configure. On a container platform, that file needs a persistent volume, or you must use an external offset store; otherwise a restart loses the offset and triggers a re-snapshot or replay.
9.3 Pricing and tier considerations
Prices change, so use the Azure pricing calculator for exact figures. What drives cost:
- Event Hubs: Standard bills by throughput units plus ingress events; Premium bills by processing units; Dedicated by capacity units. The Kafka endpoint is not available on Basic. For a claims system with modest event volume, Standard with a small number of throughput units is normally enough.
- PostgreSQL Flexible Server: compute tier (Burstable, General Purpose, Memory Optimized), storage, and backup. Logical decoding adds WAL and CPU overhead, and unconsumed slots increase storage use, so size storage with headroom.
- Azure SQL Database: pick a tier that supports CDC (DTU S3+ or vCore). CDC capture and cleanup jobs consume some of the database’s resources.
- Container Apps: a single always-on replica for the relay, so there is a small constant cost; it does not scale to zero because the relay must run continuously.
- Hidden cost: engineer time for operating connectors. Compare this against a polling worker before choosing.
9.4 Reference architecture (text)
- The Angular app calls the Claims API (ASP.NET Core on .NET 10) through Azure API Management or Front Door.
- The Claims API writes to
claimsandoutboxin one transaction on Azure Database for PostgreSQL Flexible Server (private endpoint). - Debezium (Kafka Connect or Debezium Server) on Azure Container Apps reads the WAL from the
claims_outbox_slotreplication slot and applies the Outbox Event Router. - Events are published to the
claims.claim.eventsevent hub (Kafka endpoint), keyed by claim id. - Fraud, Notification, and Read-Model services consume in their own consumer groups, apply idempotent handlers, and update their own stores (the read model in Azure SQL or PostgreSQL, served to the Angular dashboard).
- Application Insights and Azure Monitor collect traces and metrics. Alerts fire on: replication slot retained WAL size, connector task state, Event Hubs consumer lag, and outbox-to-consumer end-to-end latency.
- Key Vault holds secrets, and managed identities are used wherever the runtime supports them.
Teaching guide for my team
10.1 Explain to a beginner in 2 minutes
“When our app saves a claim, it also saves a small note in an outbox table in the same database transaction. The database keeps a diary of every change it makes. Debezium reads that diary, spots the new note, and sends it to Kafka so other services (fraud, email) can react. We do not keep asking the database ‘anything new?’ every few seconds; we just read the diary as it is written. That makes it fast and light. The one catch: Debezium might occasionally send the same note twice after a restart, so the receiving service must be able to ignore duplicates.”
10.2 Explain to an intermediate developer in 5 minutes
- Recap the dual-write problem and the outbox (Day 12): two writes in one local transaction.
- Show why polling has a latency floor and a constant query cost.
- Explain the log: WAL in PostgreSQL through logical decoding, SQL Server through CDC change tables. Debezium reads it and remembers its position (offset).
- Show the connector JSON and the Event Router:
aggregate_typepicks the topic,aggregate_idpicks the partition key,payloadbecomes the message. - Explain delivery: at-least-once, so the consumer needs an idempotent handler keyed on the event id.
- Cover the operational risks: replication slot lag and disk usage in PostgreSQL, CDC retention in SQL Server. Show where the alerts live.
- Finish with the decision rule: low volume, use polling; high volume or low latency, use log tailing.
10.3 Hands-on exercise
Goal: see a claim event flow from the database to a consumer without any publishing code in the API.
- Start PostgreSQL (with
wal_level=logical), Kafka, and Kafka Connect with Debezium using Docker Compose (the Debezium tutorial images work). - Create the
claimsandoutboxtables from section 5. - Register the connector from section 6.1 (point it to your local hosts, drop the SSL and secrets settings).
- Run the
ClaimSubmittercode to insert three claims. - Run a console consumer (
kafka-console-consumerwith--property print.key=true --property print.headers=true) onclaims.claim.events. - Stop Kafka Connect, submit two more claims, restart Kafka Connect.
- Query
SELECT slot_name, active, confirmed_flush_lsn FROM pg_replication_slots;before, during, and after the outage.
Expected outcome: the consumer shows one message per claim with the claim id as key and the event id in the headers; the two claims submitted during the outage appear after restart, in order, with no loss; the slot shows active = false during the outage, and you can explain why WAL retention grows while it is inactive. Bonus: stop the consumer mid-way, restart, and show that a duplicate is skipped by the idempotent handler.
10.4 Interview-style questions
- Q: Why is log tailing better than a polling publisher for the outbox, and when is it not? A: It removes poll latency and query load because it reads the database log sequentially and preserves commit order. It is not better at low volume or when the platform does not allow logical decoding or CDC, because it adds components (connector, slot, offsets) that the team must run and monitor.
- Q: What is the biggest operational risk of a PostgreSQL replication slot, and how do you mitigate it? A: If the consumer of the slot stops, PostgreSQL retains WAL until it is confirmed, which can fill the disk and stop the database. Mitigate with alerts on slot retained WAL,
max_slot_wal_keep_size, removing unused slots, and a runbook to restart or re-snapshot. - Q: Does log tailing give exactly-once delivery? A: No. It gives at-least-once delivery, since after a restart it can resend events from the last committed offset. Consumers must be idempotent, for example by storing processed event ids with a unique key.
Mastery checklist
- I can explain the dual-write problem and why the outbox plus a relay fixes it.
- I can describe how a PostgreSQL logical replication slot and SQL Server CDC differ in how Debezium gets the changes.
- I can configure the Debezium Outbox Event Router (topic routing, message key, payload, event id header).
- I can choose between log tailing and a polling publisher for a given workload and justify it.
- I know the failure modes (slot lag and WAL growth, CDC retention, duplicates, schema changes) and how to alert on each.
- I have written an idempotent consumer for the events.
- I can host the relay on Azure (Flexible Server settings, Event Hubs Kafka endpoint on Standard or higher, Container Apps single replica) and secure it.
- I can write an ADR that compares this with polling and states the operational cost honestly.
Key takeaway
Log tailing turns the database’s own transaction log into a fast, ordered, low-overhead event feed for your outbox, but it delivers at-least-once and it must be monitored, because a stalled reader can fill your database’s disk.
