Overview

A data mesh is an organisational model before it is a technology choice. It rests on four principles: domain-oriented data ownership, data products as first-class citizens, self-serve data infrastructure, and federated computational governance. In a banking context, the practical translation is this: the payments team owns the data that describes payments, the accounts team owns account data, and neither team should need a central integration team to access the other’s data in near-real-time.

The traditional central data warehouse fails at the speed that SAMA Open Banking requires at the integration layer. By the time the ETL job lands payment data in the warehouse, the downstream consumer that needs to power a real-time account aggregation view for an open banking API has already missed its SLA. Batch pipelines with overnight runs are not a viable architecture when the regulatory expectation is that an account holder can see all their transactions across institutions within seconds of a payment settling.

An event-driven data mesh addresses this by making the domain’s real-time event stream — a Kafka topic — its output data product. Consumers subscribe to that topic as they would subscribe to an API. The domain team owns the producer, the schema, and the SLA. The contract is not a SQL table schema that changes without notice; it is an Avro schema versioned and enforced by a Schema Registry that rejects breaking changes at publish time.

DimensionCentral integration hubData mesh
Data ownershipCentral integration teamDomain team
Schema governanceCentral data dictionary (often stale)Schema Registry per domain, versioned Avro
Cross-domain querySQL join in the data warehouseKafka Streams join or Flink SQL over event streams
Consumer discoveryData catalogue maintained by integration teamSchema Registry subjects + topic catalogue
SAMA compliance complexityOne perimeter to enforcePer-domain enforcement; federated governance required
Change velocityChange request to integration team; weeksDomain team publishes new schema version; hours

Data Products as Kafka Topics

In a data mesh, the output port of a domain is a Kafka topic with a versioned Avro schema registered in the Schema Registry. The domain team owns the producer application, owns the Avro schema definition, and owns the SLA commitment: message freshness within N seconds of the source event, partition count sufficient for the expected consumer throughput, and a schema evolution policy that does not break registered consumers without a migration window.

Consumers subscribe to that topic as data product consumers. They do not call an internal API that the domain team has to build separately; they read from the Kafka topic with a consumer group, and the Schema Registry enforces that the schema they compiled against is compatible with what the producer is publishing. The contract is the schema; the Schema Registry is the contract enforcer.

A data mesh does not mean every team manages its own Kafka cluster

The infrastructure is shared. What is distributed is ownership, schema, and SLA accountability. A single Confluent or Red Hat AMQ Streams cluster runs on shared OpenShift infrastructure managed by the platform team. Domain teams own their topics, their producer deployments, and their Schema Registry subjects — not the brokers, the ZooKeeper replacement (KRaft), or the network infrastructure. This is the same separation of concerns as application teams owning their microservices without owning the Kubernetes nodes they run on.

Banking Domain Decomposition

Domain decomposition in a banking data mesh follows the bounded context boundaries that the organisation already has — or should have. Four domains cover the core integration surface:

  • Payments domain. Owns events that describe the lifecycle of a payment instruction: PaymentInitiated (pacs.008 equivalent), PaymentSettled, PaymentFailed, PaymentReturned. Also owns settlement batch events. The payments domain does not own account balance state — that is a downstream effect it publishes for the accounts domain to consume.
  • Accounts domain. Owns balance events (BalanceUpdated, triggered by settled payments and other debits/credits), statement events (StatementReady, StatementLineAdded), and account lifecycle events (AccountOpened, AccountClosed). The accounts domain consumes payment events to update balances; it does not own the payment data.
  • Customer domain. Owns KYC events (KYCInitiated, KYCCompleted, KYCRevoked), onboarding events (OnboardingStarted, OnboardingCompleted), and consent events (ConsentGranted, ConsentRevoked). Customer domain events carry the highest PDPL classification weight — every consumer ACL in this domain requires a data classification justification.
  • Products domain. Owns rate change events (RateChanged, ProductTermsUpdated), product availability events, and fee schedule updates. Products domain events are relatively low volume but high fan-out — many downstream consumers need rate changes immediately for pricing recalculations.

The decomposition is not arbitrary. It maps to the teams that have the business knowledge to maintain the schema and the SLA. A data product owned by a team that does not understand the business rules of the data it publishes will drift into inconsistency the first time a regulatory requirement changes.

CDC as the Source of Domain Events

The practical source of most domain events in a brownfield banking environment is CDC: Debezium reading the payments database change log and transforming row-level changes into domain events. The Debezium connector captures every INSERT, UPDATE, and DELETE on the payments table and publishes it to a raw CDC topic. A Kafka Streams transformation application consumes the raw CDC topic and produces the domain event topic.

The transformation step is critical. A raw Debezium event for a payment row update is a database record change — it contains every column, including internal state columns that have no meaning outside the database. The domain event is a business event: it carries only the fields that are part of the domain contract, expressed in domain language, with a schema that is stable across database refactors.

Three transformation decisions matter in production:

  • Filter internal state changes from domain events. A status column that transitions through intermediate states (PENDING → PROCESSING → SETTLED) should not produce a domain event on every transition unless the intermediate states are meaningful to consumers. Only the transitions that represent a business state change the consuming domain cares about should become domain events.
  • Set the Kafka message key to the business entity key, not the database primary key. For a payment, the key is the end-to-end transaction reference (UETR in SWIFT, payment reference in SARIE), not an internal auto-increment ID. This ensures that all events for the same payment land in the same partition — preserving causal order for consumers that need it.
  • Tombstone events for deletes. When a row is deleted from the source, publish a tombstone (null value, entity key as the Kafka message key) to the domain event topic. Consumers that maintain a local materialised view can use the tombstone to delete their local copy. Without tombstones, deleted records linger in downstream stores indefinitely.

Schema Contracts Between Domains

The Avro schema registered in the Schema Registry is the API contract between a data product producer and its consumers. It is not documentation; it is an enforceable constraint. The Schema Registry configured with BACKWARD compatibility mode will reject any schema publication that breaks existing consumers: adding a field without a default, removing a required field, changing a field type. A breaking change cannot be pushed to production without first creating a new schema subject version and allowing a consumer migration window.

PaymentSettled.avscjson
{
  "type": "record",
  "name": "PaymentSettled",
  "namespace": "info.wbadawi.payments.v1",
  "doc": "A payment has settled successfully. Data classification: CONFIDENTIAL. PDPL category: financial-transaction. Owner: payments-domain-team@saib.com.sa. SLA: published within 5s of settlement confirmation.",
  "fields": [
    {
      "name": "uetr",
      "type": "string",
      "doc": "Unique End-to-end Transaction Reference (ISO 20022). Not PII. Classification: INTERNAL."
    },
    {
      "name": "amount",
      "type": { "type": "bytes", "logicalType": "decimal", "precision": 18, "scale": 2 },
      "doc": "Settlement amount. Classification: CONFIDENTIAL."
    },
    {
      "name": "currency",
      "type": "string",
      "doc": "ISO 4217 currency code. Classification: INTERNAL."
    },
    {
      "name": "debtorIban",
      "type": ["null", "string"],
      "default": null,
      "doc": "Debtor IBAN. PII. Classification: RESTRICTED. Access requires consumer ACL approval."
    },
    {
      "name": "creditorIban",
      "type": ["null", "string"],
      "default": null,
      "doc": "Creditor IBAN. PII. Classification: RESTRICTED. Access requires consumer ACL approval."
    },
    {
      "name": "settlementTimestamp",
      "type": { "type": "long", "logicalType": "timestamp-millis" },
      "doc": "UTC epoch milliseconds of settlement confirmation. Classification: INTERNAL."
    },
    {
      "name": "channel",
      "type": { "type": "enum", "name": "Channel",
                "symbols": ["SARIE", "SWIFT", "RTGS", "IPS", "INTERNAL"] },
      "doc": "Payment rail used for settlement. Classification: INTERNAL."
    }
  ]
}

The schema evolution policy for data products in the mesh:

  • Adding a field with a default value is BACKWARD-compatible. Existing consumers deserialise without the field; new consumers read the default. This is the safe change path.
  • Removing a field requires a new subject version (e.g., payments.settled.v2), a documented consumer migration window (minimum 30 days), deprecation notices on the v1 subject, and a committed decommission date. No silent removals.
  • Changing a field type is always a major version bump. No exceptions. Type changes break existing consumers at deserialisation time in ways that are silent unless consumers log schema errors.

Access Control on Data Products

Kafka ACLs implement the consumer authorisation model. Every topic is owned by its domain team (producer ACL), and consumer access is granted per consumer group after a documented access request. In a SAMA-regulated environment, the access request process must record the data classification, the consuming team, the consuming system, and the business justification. The ACL grant is the technical enforcement; the process record is the audit trail.

kafka-acl-setup.shbash
#!/bin/bash
# Domain ownership: payments-domain-service is the only authorised producer
kafka-acls.sh --bootstrap-server kafka.internal:9093 \
  --command-config /etc/kafka/admin.properties \
  --add \
  --allow-principal User:payments-domain-service \
  --operation Write \
  --operation Describe \
  --topic payments.settled.v1

# Consumer authorisation: open-banking-api service, read-only
kafka-acls.sh --bootstrap-server kafka.internal:9093 \
  --command-config /etc/kafka/admin.properties \
  --add \
  --allow-principal User:open-banking-api-service \
  --operation Read \
  --operation Describe \
  --topic payments.settled.v1 \
  --group open-banking-payments-cg

# Consumer authorisation: fraud-engine, read-only
kafka-acls.sh --bootstrap-server kafka.internal:9093 \
  --command-config /etc/kafka/admin.properties \
  --add \
  --allow-principal User:fraud-engine-service \
  --operation Read \
  --operation Describe \
  --topic payments.settled.v1 \
  --group fraud-engine-payments-cg

# Schema Registry: producer can write schema subjects; consumers can read only
# Enforced via Schema Registry RBAC (Confluent) or sr-acl plugin

# Verify ACLs on the topic
kafka-acls.sh --bootstrap-server kafka.internal:9093 \
  --command-config /etc/kafka/admin.properties \
  --list \
  --topic payments.settled.v1
Broadcasting PII on a Kafka topic without PDPL-compliant ACLs is a data breach at scale

A Kafka topic that carries customer name, IBAN, or Saudi National ID without consumer ACL enforcement is not a data product — it is a data breach that runs continuously. Every consumer group that subscribes to that topic receives PII without a legal basis, without a data processing record, and without the data subject’s informed consent as required by the Saudi Personal Data Protection Law. The scale multiplier of Kafka makes this worse, not better: a misconfigured ACL exposes PII to every consumer that can reach the broker. Enforce consumer ACLs for every topic that carries PDPL-classified data before the topic receives its first production message.

Cross-Domain Query Patterns

The most common cross-domain query in a banking data mesh is the enriched payment view: join a PaymentSettled event with the current account balance from the accounts domain to produce a view that shows both the payment and the post-settlement balance. This cannot be a synchronous call to the accounts service — at settlement volume, that call would saturate the accounts service. It must be a stream join.

EnrichedPaymentView.java (Kafka Streams)java
// Kafka Streams topology: join payments.settled.v1 with accounts.balance.v2
// Both topics must be co-partitioned on IBAN for the join to work without repartition
StreamsBuilder builder = new StreamsBuilder();

KStream<String, PaymentSettled> payments =
    builder.stream("payments.settled.v1",
        Consumed.with(Serdes.String(), paymentSerde));

// GlobalKTable for account balance: small-ish, needs low-latency lookup
// If balance topic is large, use KTable with co-partitioned keys instead
GlobalKTable<String, AccountBalance> balances =
    builder.globalTable("accounts.balance.v2",
        Consumed.with(Serdes.String(), balanceSerde),
        Materialized.as("account-balance-store"));

KStream<String, EnrichedPayment> enriched = payments
    .join(
        balances,
        // key mapper: extract IBAN from PaymentSettled to look up in balance table
        (paymentKey, payment) -> payment.getDebtorIban(),
        // value joiner: merge payment event with current balance snapshot
        (payment, balance) -> EnrichedPayment.builder()
            .uetr(payment.getUetr())
            .settledAmount(payment.getAmount())
            .currency(payment.getCurrency())
            .settlementTimestamp(payment.getSettlementTimestamp())
            .postSettlementBalance(balance.getAvailableBalance())
            .accountCurrency(balance.getCurrency())
            .build()
    );

// Publish to a new data product topic owned by the cross-domain team
enriched.to("payments.enriched.v1",
    Produced.with(Serdes.String(), enrichedSerde));

// Also expose via interactive query (REST API over state store)
// Allows REST clients to query current balance without subscribing to Kafka
ReadOnlyKeyValueStore<String, AccountBalance> balanceStore =
    streams.store(StoreQueryParameters.fromNameAndType(
        "account-balance-store", QueryableStoreTypes.keyValueStore()));
Cross-domain joins in Kafka Streams are expensive when partition counts differ

Kafka Streams joins between a KStream and a KTable require both topics to be co-partitioned: same number of partitions, same partitioning key. If payments.settled.v1 has 12 partitions and accounts.balance.v2 has 6 partitions, Kafka Streams will silently repartition the smaller topic — creating an internal repartition topic, adding latency, and doubling broker storage for the duration of the join. Design both topics with the same partition count from the start, keyed on the same business entity key (IBAN or account number). Changing partition counts after consumers are live requires a coordinated migration that is far more painful than getting the count right initially.

For complex cross-domain aggregations that span more than two domains — for example, a regulatory report that joins payment events with account state and customer KYC status — Flink SQL is more appropriate than Kafka Streams. Flink’s SQL layer treats Kafka topics as tables and expresses multi-way joins as standard SQL, with temporal joins for handling slowly-changing dimension data (KYC status changes infrequently; payment events arrive continuously).

Federated Governance Model

Federated governance is the hardest part of a data mesh to implement in a regulated bank. The principle is that domain teams are accountable for their data products’ quality, classification, and compliance, while a central governance body sets the policies that domain teams must implement. In practice, this means the central governance team cannot manually approve every schema change, every topic creation, or every consumer ACL grant — that recreates the bottleneck you were trying to remove. Instead, governance enforces policy through tooling and automated checks.

  • Data product catalogue. Every Schema Registry subject is an entry in the data product catalogue. The catalogue records: domain owner, data classification label, consumer list, SLA, PDPL processing basis, and retention policy. This catalogue is maintained as code (a Git repository), and schema publication pipelines validate that the catalogue entry exists and is current before registering a new schema version.
  • Data quality SLA per domain. Each domain team publishes data quality metrics for their topics: freshness (time between source event and domain event publication, p99), completeness (fraction of source events that produced a domain event), and schema validity rate (fraction of messages that deserialise against the registered schema). These are Prometheus metrics. Violations alert to the domain team first, then escalate to governance if unresolved within the SLA window.
  • Central governance enforces PDPL and SAMA data classification. The governance layer sets the classification labels, validates that the schema doc fields carry the correct PDPL metadata, and audits consumer ACL grants against the classification policy. A consumer cannot be granted read access to a RESTRICTED topic without a documented data processing record in the PDPL register.
  • Domain team accountability for data product SLA. Domain teams are on-call for their data products. If payments.settled.v1 has a freshness SLA of 5 seconds p99 and the 5-second threshold is breached, the payments team is paged, not the platform team. Platform is accountable for broker availability; domain teams are accountable for producer behaviour and schema quality.

SAMA Data Residency in a Mesh

SAMA requires that all banking data remains within the Kingdom of Saudi Arabia. In a data mesh, this constraint applies at every layer: the Kafka brokers, the Schema Registry, the ZooKeeper/KRaft quorum, the consumer state stores, and the Flink job state backends must all run on infrastructure physically located within Saudi Arabia. There is no exception for development, staging, or disaster recovery environments that replicate regulated data.

  • All Kafka topics reside on infrastructure within Saudi Arabia. Red Hat AMQ Streams on OpenShift running in a Saudi data center, or Confluent Platform in a SAMA-approved cloud region. Cross-border replication (MirrorMaker 2, Confluent Cluster Linking) to a DR site outside Saudi Arabia is prohibited without explicit SAMA approval under the data residency framework.
  • Schema Registry within the regulated perimeter. The Schema Registry is not a stateless service — it stores schema versions that map to data product contracts. If the Schema Registry is hosted outside Saudi Arabia, schema subject registration and schema validation happen outside the regulatory perimeter, which is a compliance gap even if the data itself stays in-country.
  • Data lineage tracking for audit. The central governance layer must maintain a record of which domain teams produce data to which topics, which consumer groups read from which topics, and which cross-domain joins materialise data from multiple topics. This lineage record is the evidence SAMA can request to map the flow of customer data through the mesh. Without it, a data residency audit cannot be completed.

Production Checklist: Publishing a Domain’s First Data Product

  1. Write the data product specification

    Before writing any code, document in the data product catalogue: the domain, the topic name and version, the business event it represents, the Kafka message key, the SLA (freshness p99, completeness), the data classification label, the PDPL processing basis, and the planned consumer list. This document is a governance artefact. Get it reviewed by the data governance team before the Schema Registry subject is created.

  2. Draft and validate the Avro schema

    Write the Avro schema with doc fields on every field that carry the classification label and PII indicator. Validate BACKWARD compatibility against any previous version of the subject. Include at least one optional field with a default (not a union type unnecessarily) to leave room for the first expected evolution without a breaking change.

  3. Register the schema and create the topic

    Register the schema subject in the Schema Registry with the compatibility mode set to BACKWARD. Create the Kafka topic with the correct partition count (co-partitioned with the topics you plan to join against), replication factor of 3, and retention configured to match the data classification retention policy (minimum 7 years for transaction data under SAMA TRM).

  4. Set producer ACL and deploy the producer

    Grant the domain service principal Write and Describe ACLs on the topic. Deploy the Debezium connector or Kafka producer application. Verify that messages are being produced with the correct schema by consuming the first batch and checking deserialisation against the Schema Registry schema.

  5. Grant consumer ACLs per approved consumer

    For each approved consumer in the data product specification, grant Read and Describe ACLs to the consumer service principal and the consumer group name. Record the ACL grant in the PDPL data processing register if the topic carries PDPL-classified data. Do not use a wildcard consumer ACL (*) for any topic that carries RESTRICTED or CONFIDENTIAL data.

  6. Instrument and publish data quality metrics

    Add Prometheus metrics to the producer: data_product_messages_published_total, data_product_source_to_topic_latency_seconds, data_product_schema_serialisation_errors_total. Wire these to the Grafana data quality dashboard. Set alerts: freshness SLA breach pages the domain team on-call; schema serialisation errors are a P1 (a consumer that receives a message it cannot deserialise will fail silently if it lacks proper error handling).

  7. Run the governance review and open the topic to consumers

    Submit the data product for governance review: data product specification, schema Registry subject registration confirmation, ACL configuration, data quality dashboard link, and PDPL register entry. Once governance signs off, announce the data product in the internal data product catalogue and notify approved consumers that the topic is production-ready.