Manikandan — Manikandan
Microservices

Day 13: Transaction Log Tailing

ManikandanManikandan
23 min read·Updated Aug 22, 2022

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 (ClaimSubmitted before ClaimApproved).
  • Fast: the fraud-check service should see ClaimSubmitted in 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 ProcessedAt compete 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:

  1. Event delivery latency tracks the poll interval. Traces show a consistent gap equal to the polling period between ClaimSubmitted committed and ClaimSubmitted consumed.
  2. Database CPU or query stats show the outbox query near the top of pg_stat_statements or SQL Server Query Store by execution count, even though it returns zero rows most of the time.
  3. Lock waits or deadlocks on the outbox table during peak intake (for example after a storm or catastrophe event that triggers many claims).
  4. Outbox table is growing because cleanup jobs cannot keep up, and the Unprocessed index gets bigger and slower.
  5. Out-of-order events on the consumer side after you scaled the relay to two instances.
  6. 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.
  7. 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:

  1. Your service writes the business row and an outbox row in one transaction.
  2. A log reader (Debezium) reads the database log and turns each new outbox row into a message.
  3. 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.

-- PostgreSQL
CREATE 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 transaction
using 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 outbox with aggregate_type = 'claim' go to the topic claims.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 payload JSON, and id is 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.Kafka
using 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 seconds
import { 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 systems

It 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_id as the message key. Never key by aggregate_type alone, or one hot partition receives everything.
  • Batching: tune max.batch.size and max.queue.size on the connector when snapshots or bursts cause memory pressure.
  • Keep the outbox narrow: keep the payload reasonably 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 modeWhat happensPrevention / 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 restartThe consumer sees the same event again.Idempotent consumer keyed on the event id header (Day 17).
Schema change on the outbox tableConnector 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 databaseSlot 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 unavailableConnector 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 debezium database user with the least privileges: REPLICATION plus SELECT on the outbox table (PostgreSQL), or db_datareader on the CDC schema (SQL Server).
  • Use TLS from connector to database (sslmode=require or 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

  1. Capturing business tables instead of the outbox. This publishes your internal schema and creates hidden coupling. Capture only the outbox.
  2. Forgetting slot monitoring (PostgreSQL) and discovering the problem when the disk is full.
  3. Assuming exactly-once. It is at-least-once; the consumer must be idempotent.
  4. Wrong message key, breaking per-claim ordering.
  5. Running two connectors with the same slot name or restarting with a new slot name and silently skipping or replaying history.
  6. 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 (or never in older versions) so old, already-published rows are not published again. Check the option names for your Debezium version.
  7. 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

CriterionTransaction Log Tailing (Debezium/CDC)Polling Publisher (Day 14)Dual write (anti-pattern)Managed change feed (e.g. Cosmos DB change feed)
Delivery guaranteeAt-least-once, ordered per keyAt-least-once, ordering needs careNone (partial failures)At-least-once, ordered per partition key
LatencySub-second to a few secondsPoll intervalLowest but unsafeLow
Load on source DBVery low (sequential log read)Continuous queries and updatesNone extraHandled by the platform
Operational complexityHigh (connector, slots, offsets, broker)Low (a .NET worker)LowestMedium
Database supportNeeds CDC/logical decoding supportAny databaseAnyOnly that database
Risk of taking down the DBYes if slot lag ignoredLowData inconsistency insteadLow
Best forHigh volume, low latency, several consumersModest volume, restricted platformsNever in productionCosmos 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, and ClaimRejected events 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 outbox table from the Claims PostgreSQL database with Debezium (logical decoding, pgoutput), route with the Outbox Event Router to topics claims.<aggregate_type>.events keyed by aggregate_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

RoleAzure serviceNotes
Source database (PostgreSQL)Azure Database for PostgreSQL - Flexible ServerSet 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 InstanceSupports 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.
BrokerAzure 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 runtimeAzure Container Apps or AKSRun Kafka Connect with the Debezium connector, or Debezium Server (which has an Azure Event Hubs sink). One replica only for a given slot.
SecretsAzure Key VaultDatabase and broker credentials, referenced through managed identity where the runtime supports it.
MonitoringAzure Monitor, Log Analytics, Application InsightsConnector logs, database metrics, Event Hubs metrics, alerts.
ConsumersAzure Container Apps / AKS / Azure Functions (Event Hubs trigger)Idempotent handlers.

9.2 How to configure

  1. PostgreSQL Flexible Server: in Server parameters set wal_level to logical, review max_replication_slots and max_wal_senders, restart the server. Create the debezium role 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 the outbox table and, if you prefer, the publication in advance.
  2. Event Hubs: create a namespace on Standard or higher. In Event Hubs, a Kafka topic is an event hub, so pre-create claims.claim.events and 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:9093 with SASL_SSL, mechanism PLAIN, username $ConnectionString, and the connection string (or a Microsoft Entra ID / OAuth setup) as the password.
  3. 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.
  4. Networking: use private endpoints for the database and Event Hubs, or VNet integration and firewall rules, so the relay reaches both privately.
  5. Consumers: use consumer groups per service (fraud-service, notification-service).

Debezium Server sink example (properties, simplified):

debezium.sink.type=eventhubs
debezium.sink.eventhubs.connectionstring=${EVENTHUBS_CONNECTION_STRING}
debezium.sink.eventhubs.hubname=claims.claim.events
debezium.source.connector.class=io.debezium.connector.postgresql.PostgresConnector
debezium.source.plugin.name=pgoutput
debezium.source.table.include.list=public.outbox
debezium.source.topic.prefix=claims-db
debezium.source.offset.storage.file.filename=/data/offsets.dat

Debezium 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)

  1. The Angular app calls the Claims API (ASP.NET Core on .NET 10) through Azure API Management or Front Door.
  2. The Claims API writes to claims and outbox in one transaction on Azure Database for PostgreSQL Flexible Server (private endpoint).
  3. Debezium (Kafka Connect or Debezium Server) on Azure Container Apps reads the WAL from the claims_outbox_slot replication slot and applies the Outbox Event Router.
  4. Events are published to the claims.claim.events event hub (Kafka endpoint), keyed by claim id.
  5. 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).
  6. 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.
  7. 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

  1. Recap the dual-write problem and the outbox (Day 12): two writes in one local transaction.
  2. Show why polling has a latency floor and a constant query cost.
  3. Explain the log: WAL in PostgreSQL through logical decoding, SQL Server through CDC change tables. Debezium reads it and remembers its position (offset).
  4. Show the connector JSON and the Event Router: aggregate_type picks the topic, aggregate_id picks the partition key, payload becomes the message.
  5. Explain delivery: at-least-once, so the consumer needs an idempotent handler keyed on the event id.
  6. Cover the operational risks: replication slot lag and disk usage in PostgreSQL, CDC retention in SQL Server. Show where the alerts live.
  7. 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.

  1. Start PostgreSQL (with wal_level=logical), Kafka, and Kafka Connect with Debezium using Docker Compose (the Debezium tutorial images work).
  2. Create the claims and outbox tables from section 5.
  3. Register the connector from section 6.1 (point it to your local hosts, drop the SSL and secrets settings).
  4. Run the ClaimSubmitter code to insert three claims.
  5. Run a console consumer (kafka-console-consumer with --property print.key=true --property print.headers=true) on claims.claim.events.
  6. Stop Kafka Connect, submit two more claims, restart Kafka Connect.
  7. 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

  1. 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.
  2. 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.
  3. 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.

Interactive Architectural Roadmaps

Explore Complete Roadmaps & Pattern Checklists

Track your learning with interactive checklists for all 23 Gang of Four patterns and modern Microservice architecture patterns.

Share:
Back to Blog

Related Posts

View All Posts
Microservices

Day 12: Transactional Outbox

The Transactional Outbox pattern solves the "dual-write" problem: a service must update its database and publish an event, but a database and a message broker cannot share one ACID transaction.

Manikandan
Manikandan·16 min read
Microservices

Day 53: Client-Side UI Composition (Micro-Frontends)

Client-Side UI Composition splits one large single-page application into several smaller front-end applications ("micro-frontends") that are built, tested, and deployed independently by different teams.

Manikandan
Manikandan·17 min read