Free Kafka Design Patterns Interview Practice Questions

Try 24 original Kafka design-pattern questions with explained answers, record traces and a workflow diagram. Practise architecture reasoning before an interview.

These 24 original IT Mastery practice questions cover eight Kafka design and integration topics. The set contains 16 single-answer questions and eight Select TWO questions. Read the requested answer count before responding.

This is independent professional-skills and interview preparation, not an official certification exam. These are not official exam questions, copied live-exam content or exam dumps. The sample length, topic weights and interaction mix are editorial choices; there is no official time limit or passing score for this track.

Try each question before revealing its explanation. For a correct guess, explain the deciding evidence in your own words and revisit it later. Kafka-specific API examples use the stated Kafka 4.2 reference version.

Practice-set coverage

DomainIT Mastery practice weightQuestions in this set
Event Modeling and Contract Boundaries10%2
Channels Partitioning and Routing Patterns15%4
Delivery Guarantees Consistency and Recovery15%4
Stateful Stream Processing and Temporal Patterns15%4
Enterprise Integration and Business Workflows15%4
Evolution Security and Data Governance10%2
Testing Observability and Controlled Operation10%2
Architecture Tradeoffs and Integrated Design Cases10%2

Practice questions

Questions 1-24

Question 1

Topic: Event modeling and contract boundaries

An order service currently publishes the following record when it wants a shipping label created.

Observed behavior:

topic: order.integration
key: O-417
type: ShipmentRequested
messageId: M-91
groups: fulfillment-primary, fulfillment-recovery, analytics
result: two labels created for O-417

Contract constraints:

  • Fulfillment owns the decision to create or reject a shipment.
  • A standby worker must take over after a rebalance, not process independently.
  • Requests must remain durable while fulfillment is unavailable.
  • Analytics needs facts distinguishing creation from business rejection.
  • Eliminating independent primary/recovery delivery is in scope; making the external label-creation side effect retry-safe remains a fulfillment responsibility.

Which bounded redesign best supports these constraints?

Options:

  • A. Call fulfillment synchronously with CreateShipment; publish ShipmentCreated or ShipmentRejected after the response.

  • B. Send a keyed CreateShipment command to separate primary and recovery groups; retain the first outcome received.

  • C. Publish a keyed ShipmentRequested event to one fulfillment group; use commits and dead-letter records as business outcomes.

  • D. Send a keyed CreateShipment command to one fulfillment group; publish ShipmentCreated or ShipmentRejected events.

Best answer: D

Explanation: Despite its past-tense name, ShipmentRequested functions as an instruction: the order service asks a specific business owner to perform work that may be rejected. A targeted command therefore expresses the contract more accurately. Members of one ordinary consumer group share partition ownership, whereas separate groups each receive the record independently; using the same key does not coordinate work across groups.

A retained command remains available during a fulfillment outage. After making the business decision, fulfillment publishes an event recording what actually happened. Kafka commits and dead-letter records describe processing state, not whether the business accepted or rejected the shipment request.

A crash after label creation but before the consumed offset is committed can still cause redelivery. Fulfillment must make that external effect retry-safe, such as by using the stable messageId for idempotency or durable deduplication; a single consumer group alone does not provide that guarantee.

Why each option fits or fails:

A. Synchronous command cannot durably await processing while fulfillment is unavailable, despite producing suitable result events afterward.

B. Separate recovery group gives both groups independent delivery, so the matching key does not prevent duplicate label creation.

C. Commits as outcomes confuses successful record processing with shipment creation and technical failure with business rejection.

D. The targeted work may be rejected by fulfillment, while one group removes deliberate independent fan-out and outcome events record business facts. The group does not by itself make external label creation exactly once.


Question 2

Topic: Event modeling and contract boundaries

A producer publishes OrderCreated records by reflectively serializing its internal domain class. Renaming the internal member customer_id to buyer_id changed the serialized field name, causing independently deployed consumers to reject new records.

The topic contract must remain stable across internal refactors while permitting compatible schema evolution. Which TWO changes should the team make?

Options:

  • A. Regenerate the topic contract from current producer classes during each build.

  • B. Coordinate all producer and consumer deployments when internal model names change.

  • C. Share producer domain classes as the canonical model for every consumer.

  • D. Establish a versioned topic contract and explicitly map domain objects to it.

  • E. Test each writer and reader against all supported contract versions in CI.

Correct answers: D and E

Explanation: A transport schema is an integration contract, not a serialization of the producer’s internal model. A separately owned transport model, combined with explicit mapping, allows the producer to rename customer_id internally while continuing to publish the contracted field. Version-aware writer and reader tests then verify that proposed schema changes remain compatible with supported deployments.

Sharing domain classes may reuse code, but it also couples consumers to the producer’s model and release cycle. Regenerating the contract from domain classes preserves the original failure mechanism, while coordinated deployment abandons the requirement for independent evolution.

Why each option fits or fails:

A. Generated from internals allows another class refactor to change the serialized contract unintentionally.

B. Coordinated deployments can manage one migration but do not support independently released producers and consumers.

C. Shared domain classes create compile-time and deployment coupling rather than a stable transport boundary.

D. An explicit boundary model prevents internal member changes from altering the transport schema.

E. Compatibility tests detect transport changes that would break independently deployed readers or writers.


Question 3

Topic: Channels partitioning and routing patterns

An order processor uses four ordinary Kafka consumers. Each consumer can process at most 500 records/s across its assigned partitions. The six-partition topic is keyed by tenantId.

PartitionRateCurrent owner
P0720/sC1
P1120/sC2
P2110/sC3
P3100/sC4
P490/sC1
P580/sC2

Backlog must not grow, and the fewest consumers should be used. Ordering is required per order, not across a tenant’s orders. A tested six-partition replacement keyed by orderId has no partition above 220 records/s. The assignor distributes partition counts evenly. During cutover, producers may be paused, and publishing to the replacement begins only after the old topic is fully drained.

Which design best meets the goal?

Options:

  • A. Retain the current route and scale the consumer group to eight consumers.

  • B. Retain the current route and scale the consumer group to six consumers.

  • C. Expand the current topic to twelve partitions and run twelve consumers.

  • D. Drain the old route, then use the replacement topic with three consumers.

Best answer: D

Explanation: An ordinary consumer group assigns each partition to at most one consumer, so replicas cannot divide P0’s 720 records/s. Six consumers would isolate P0 but would not raise its 500 records/s processing limit; consumers beyond the six partitions would be idle.

Keying the replacement topic by orderId matches the required ordering scope. With three consumers, each receives two partitions, so the tested per-partition maximum gives a projected load of at most 440 records/s. Two consumers are insufficient because the total input rate is 1,220 records/s versus 1,000 records/s of combined capacity. Thus, three is the smallest supported group after rerouting.

Why each option fits or fails:

A. Eight consumers leave two instances idle because the topic has only six partitions, while P0 remains overloaded.

B. Six consumers cannot split P0, so its 720 records/s still exceeds one consumer’s capacity.

C. Twelve partitions are not proven by the aggregate rates to bring every partition below 500 records/s and would use more consumers than the proven three-consumer design. In-place expansion would also require an ordering-safe cutover because changing the partition count can remap existing keys.

D. Three consumers each receive two partitions, limiting projected load to 440 records/s while preserving per-order ordering.


Question 4

Topic: Channels partitioning and routing patterns

An order projection reads a topic keyed by customerId. The topic has eight partitions and eight ordinary consumer-group members. Each partition can sustain 1,000 records/s.

TrafficArrival rateDistribution
Customer HOT2,400 records/sOne unchanged key
Other customers3,200 records/sAbout 400 per partition

Lag grows only on HOT’s partition. Events include a monotonic customerSequence, and the projection requires per-customer ordering unless redesigned handling reassembles that sequence or relaxes the requirement.

The team is considering more partitions, a dedicated one-partition topic for HOT, or three evenly loaded HOT shard keys explicitly routed to three reserved partitions, with other customers rerouted across the remaining five partitions.

Which TWO conclusions are supported?

Options:

  • A. A dedicated one-partition HOT topic can meet the rate by removing other customer traffic.

  • B. Increasing to 24 partitions cannot remove HOT’s overload while its key remains unchanged.

  • C. Three HOT shard keys can meet the rate only with sequence reassembly or relaxed customer ordering.

  • D. Producer idempotence can distribute HOT across partitions while retaining its original customer order.

  • E. Adding 16 consumers to the current group can divide HOT’s records among idle members.

Correct answers: B and C

Explanation: Kafka preserves ordering within a partition, and one unchanged key is routed to one partition. HOT therefore sends 2,400 records/s to a resource that processes only 1,000 records/s. Isolating the key removes approximately 400 records/s of unrelated traffic but still leaves HOT above capacity. Adding partitions or consumers does not divide one partition’s work among ordinary group members.

Three evenly loaded shards provide 3,000 records/s of aggregate capacity, which exceeds HOT’s arrival rate. However, Kafka does not preserve customer-wide order across those partitions. The projection must reorder records using customerSequence, including an appropriate buffering rule, or explicitly accept relaxed ordering. Producer idempotence does not change this ordering boundary.

Why each option fits or fails:

A. Dedicated partition still receives 2,400 records/s against 1,000 records/s capacity, even after other traffic is removed.

B. HOT remains limited to one partition’s 1,000 records/s capacity despite the additional partitions.

C. Three partitions provide 3,000 records/s aggregate capacity but no longer provide customer-wide partition ordering.

D. Producer idempotence controls duplicate production within its protocol scope; it neither shards a key nor creates cross-partition order.

E. Additional consumers cannot concurrently process the same partition in an ordinary consumer group.


Question 5

Topic: Channels partitioning and routing patterns

A Kafka router rebuilds an empty projection by replaying order events. Historical replays must reproduce each event’s original destination. At acceptance, ingress persists the active ruleVersion in the immutable event. All rule versions are retained.

v6, active before 10:00: all orders -> standard
v7, active 10:00-10:10: EU orders -> eu-review
v8, active after 10:10: amount >= 10000 -> risk-review

E1: occurredAt=09:55, acceptedAt=10:05,
    region=EU, amount=12000, ruleVersion=7
Original E1 destination: eu-review
Replay starts: 11:00

The router commits its output and consumed offset in one Kafka transaction. versionAt(t) returns the rule active at wall-clock time t.

on Process(e):
    r = rules.getRequired(activeRuleVersion)  // replace this line
    emit(r.route(e))

Which replacement preserves reproducible routing?

Options:

  • A. Use rules.getRequired(versionAt(e.occurredAt)).

  • B. Use rules.getRequired(e.ruleVersion).

  • C. Use rules.getRequired(replayStartVersion).

  • D. Use rules.getRequired(versionAt(processingTime())).

Best answer: B

Explanation: Reproducible routing requires treating the applied rule version as event provenance. E1 occurred at 09:55 but was accepted at 10:05, when v7 was active, so its persisted ruleVersion is authoritative. Loading v7 during replay again routes the EU order to eu-review.

The Kafka transaction makes the output and consumed offset atomic within Kafka, but it cannot prevent semantic drift caused by selecting a different rule. Retaining immutable rule versions and failing when the recorded version is unavailable preserves deterministic replay. A current or replay-start rule is appropriate only when the business requirement explicitly calls for reinterpretation under current policy.

Why each option fits or fails:

A. Occurrence-time lookup selects v6 because 09:55 predates v7, ignoring the acceptance-time routing decision.

B. The persisted version selects v7, reproducing E1’s original eu-review destination.

C. Replay-start pinning consistently applies v8, but consistency within one replay does not reproduce each event’s historical decision.

D. Processing-time lookup selects v8 at 11:00, changing the destination to risk-review.


Question 6

Topic: Channels partitioning and routing patterns

A quote service aggregates replies keyed by request_id in one Kafka partition. It currently closes after 500 ms without a reply.

  • Providers may be discovered and invited until 2,000 ms after request acceptance; no membership list is sent.
  • The lowest valid quote received by 2,000 ms must be selected, with no minimum response count.
  • Results must include status, response count, and responder IDs.
ProviderArrivalQuote
A300 ms$100
B700 ms$98
C1,300 ms$96
D1,900 ms$92

Which bounded design change best satisfies these requirements?

Options:

  • A. Snapshot providers at acceptance and close when all reply or at 2,000 ms; emit ALL_RECEIVED with response metadata.

  • B. Close after three distinct replies or at 2,000 ms; emit THRESHOLD_MET with response metadata.

  • C. Close only at 2,000 ms; emit DEADLINE_CLOSED with the best quote and response metadata.

  • D. Close after 500 ms of silence or at 2,000 ms; emit IDLE_CLOSED with response metadata.

Best answer: C

Explanation: A deadline completion condition fits when membership remains unknown and every response arriving through a fixed cutoff matters. The current inactivity window is sliding: after provider B replies at 700 ms, it closes at 1,200 ms and excludes C and D. A threshold of three closes at 1,300 ms and misses D’s lower quote. Because providers can still be discovered during the collection period, an acceptance-time roster also changes the required membership semantics. At 2,000 ms, the service should select the best collected quote and emit DEADLINE_CLOSED, the response count, and responder IDs. These fields expose the bounded result without claiming that every possible provider responded.

Why each option fits or fails:

A. Snapshot membership excludes providers discovered after acceptance and therefore changes the defined collection population.

B. Threshold completion closes after C responds, so D’s timely lower quote is not considered.

C. The fixed deadline includes every timely response while explicitly avoiding an unsupported claim that all unknown members replied.

D. Inactivity completion closes at 1,200 ms, excluding valid replies that arrive before the business deadline.


Question 7

Topic: Delivery guarantees consistency and recovery

A telemetry gateway groups independent device events only to improve producer throughput. Consumers do not deduplicate event IDs.

Constraints:

  • Successful events must advance independently.
  • Permanent failures must be quarantined for correction.
  • Asynchronous sending must remain enabled.

The producer uses acks=all and enable.idempotence=true but no transaction:

for (Event event : batch) {
    producer.send(toRecord(event), (metadata, error) ->
        results.put(event.id(), error));
}
producer.flush();
if (results.values().stream().anyMatch(Objects::nonNull))
    retry(batch);

Callbacks report success for e17, RecordTooLargeException for e18, and success for e19. Which bounded change best preserves the required business semantics?

Options:

  • A. Retain batch-level failure handling, resend the complete list, and rely on producer idempotence.

  • B. Track callbacks by event ID, finalize successes, and quarantine permanent failures for correction.

  • C. Wrap each list in one transaction, abort partial lists, and require consumers to read committed records.

  • D. Send each event synchronously, finalize each acknowledgment, and quarantine a failure before continuing.

Best answer: B

Explanation: Producer batching does not make independent sends an atomic business unit. flush() waits for outstanding sends to complete, but each callback still reports its own outcome. Here, e17 and e19 were acknowledged, while e18 permanently failed. Retrying the entire list would reissue successful business events.

Producer idempotence protects against duplicates from supported protocol-level retries. It does not deduplicate new send() calls made by application-level batch retry logic. Recording completion per event allows the gateway to finalize successful events and quarantine only e18, while preserving asynchronous throughput.

A transaction would create an all-or-none boundary that incorrectly couples otherwise independent events.

Why each option fits or fails:

A. Complete-list retry can duplicate acknowledged events because producer idempotence does not deduplicate arbitrary application resubmissions.

B. Per-event callback handling preserves independent outcomes and avoids reissuing acknowledged events while retaining asynchronous batching.

C. Transactional list delays successful independent events because one permanent failure aborts the entire list.

D. Synchronous sending preserves per-event outcomes but violates the requirement to retain asynchronous throughput.


Question 8

Topic: Delivery guarantees consistency and recovery

A service uses an ordinary Kafka consumer group and an external database:

  • Record 809 is the last completed record for partition p3, and its committed consumer position is 810.
  • Worker W1 owns p3 with lease epoch 17 and submits the mutation for the record at offset 810.
  • A rebalance assigns p3 to W2 and atomically advances its lease epoch to 18.
  • W1 requests cancellation, but the database documents cancellation as best effort after acceptance.
  • W2 resumes from the record at offset 810. Each mutation carries its lease epoch, and the database supports transactional conditional updates.

Which TWO conclusions are supported?

Options:

  • A. Rely on single-owner partition assignment to reject W1’s request after reassignment.

  • B. Use a revocation-time commit of position 811 as the correctness mechanism to suppress W2’s replay of record 810.

  • C. Use static group membership to make W1’s accepted mutation lose validity automatically.

  • D. Require each database mutation to atomically match the current partition lease epoch.

  • E. Treat cancellation as advisory because W1’s accepted mutation can still complete.

Correct answers: D and E

Explanation: Kafka partition assignment controls which consumer may fetch and process new records, but it cannot revoke work already submitted to an external system. A monotonically increasing ownership epoch provides a fencing token. The database must compare the mutation’s epoch with the authoritative current epoch and apply the mutation in the same transaction as that comparison. W1’s delayed epoch-17 request is then rejected after epoch 18 is activated.

Best-effort cancellation may reduce unnecessary work, but it is not a correctness boundary. Offset management also remains separate: committing before confirmed completion can lose the mutation, while delaying the commit permits replay but does not fence stale external work.

Why each option fits or fails:

A. Group exclusivity governs Kafka partition delivery, not the external database’s authorization to complete prior work.

B. Prematurely committing position 811 suppresses replay of record 810 but can lose its mutation if W1’s accepted request is later rejected by epoch fencing; it also does not cancel that request.

C. Static membership may reduce rebalances but does not invalidate external requests when ownership eventually changes.

D. An atomic epoch comparison rejects W1’s epoch-17 mutation after epoch 18 becomes authoritative.

E. The cancellation contract does not revoke a mutation that the external database has already accepted.


Question 9

Topic: Delivery guarantees consistency and recovery

An order processor consumes input partition 0 at offset 40 and writes a result to the same Kafka cluster. Before processing, the group’s committed next offset is 40.

producer.beginTransaction();
producer.send(resultRecord);
producer.sendOffsetsToTransaction(nextOffsets,
                                  consumer.groupMetadata());
producer.commitTransaction();

Failure trace:

  • Attempt 1 writes O1 and sends next offset 41, then crashes before commitTransaction().
  • Its replacement initializes the same transactional instance, aborting the unfinished attempt.
  • Attempt 2 rereads offset 40, writes O2, sends offset 41, and commits.
  • The result log physically contains O1 and O2.

Business consumers must act only on committed results. A forensic consumer must retain visibility into physical records. Which design best satisfies these requirements?

Options:

  • A. Retain the transaction; use read_uncommitted with input-offset deduplication for business.

  • B. Retain the transaction; use read_committed for business and read_uncommitted for forensics.

  • C. Commit the result transaction first; then commit the input offset separately.

  • D. Commit the input offset first; then commit the result transaction separately.

Best answer: B

Explanation: Kafka transactions atomically commit produced records and offsets supplied through sendOffsetsToTransaction. The crash occurs before that commit, so attempt 1’s output and offset update are both aborted. The output record may remain physically present in the partition with transactional metadata; abortion does not mean the bytes were never written.

A read_committed consumer skips aborted records and therefore processes only O2. A read_uncommitted forensic consumer can observe both physical attempts. After the abort, the input offset remains eligible for replay until attempt 2 commits O2 and next offset 41. This guarantee covers Kafka records and consumer offsets, not external business effects.

Why each option fits or fails:

A. Offset deduplication can suppress repeats but may let an aborted record trigger a business effect before its transaction outcome is known.

B. The commit atomically exposes O2 and offset 41, while read_committed excludes the physically retained aborted O1.

C. Output-first commit permits a crash after output commit but before offset commit, causing another committed output on replay.

D. Offset-first commit permits a crash after offset commit but before output commit, permanently skipping the result.


Question 10

Topic: Delivery guarantees consistency and recovery

An order service fails from cluster A to asynchronously mirrored cluster B. No transaction spans both clusters or the fulfillment API. The fulfillment API enforces durable idempotency by eventId, with keys retained for the recovery period, and B uses auto.offset.reset=latest.

Partition 3 at failure:

FactPosition
Fulfillment completed on Athrough 8,080
Latest offset translationA 8,060 -> B 12,500
Last record mirrored to BA 8,090 -> B 12,530
Last record acknowledged on A8,110
Earliest retained B offset12,000

Recovery may omit at most 20 acknowledged records per partition. Which failover contract is best supported?

Options:

  • A. Resume after B offset 12,530; deduplicate by eventId; accept A records 8,091-8,110 as the bounded gap.

  • B. Resume after B offset 12,500; deduplicate by eventId; accept A records 8,091-8,110 as the bounded gap.

  • C. Resume after B offset 12,500; rely on producer idempotence; accept A records 8,091-8,110 as the bounded gap.

  • D. Resume after B offset 8,080; deduplicate by eventId; accept A records 8,091-8,110 as the bounded gap.

Best answer: B

Explanation: Kafka offsets are positions within a specific topic partition and are not portable across clusters. The durable translation maps A offset 8,060 to B offset 12,500, so recovery starts after that B record. This replays A records 8,061-8,090. The eventId registry suppresses repeated fulfillment for 8,061-8,080, while 8,081-8,090 receive their first effect. Records 8,091-8,110 were acknowledged on A but never mirrored, creating the permitted 20-record loss.

Starting at B’s mirror frontier would skip ten recoverable, unprocessed records.

Why each option fits or fails:

A. Start at mirror frontier skips mirrored A records 8,081-8,090 that have not completed fulfillment.

B. The translated checkpoint supports safe replay, eventId prevents repeated effects, and the 20 unmirrored records equal the permitted loss.

C. Producer idempotence cannot deduplicate fulfillment API calls or arbitrarily replayed business events.

D. Reuse source offsets fails because B offset 8,080 is out of range and resets to B’s latest position.


Question 11

Topic: Stateful stream processing and temporal patterns

A Kafka Streams topology has caching disabled. Both state stores start empty, and records for each key are processed in partition-offset order.

KTable<String, Order> orders = builder.table("orders");

KTable<String, Order> openTable = orders.filter(
    (key, value) -> value.status().equals("OPEN"),
    Materialized.as("open-table"));

KTable<String, Order> openViaStream = orders.toStream()
    .filter((key, value) -> value.status().equals("OPEN"))
    .toTable(Materialized.as("open-via-stream"));

The same partition supplies these updates sequentially:

offset 20: key=o-17, revision=1, status=OPEN
offset 21: key=o-17, revision=2, status=CLOSED

After both updates are processed, which TWO conclusions are correct?

Options:

  • A. open-via-stream contains revision 2 for o-17.

  • B. open-table retains revision 1 for o-17.

  • C. open-via-stream retains revision 1 for o-17.

  • D. open-via-stream has no row for o-17.

  • E. open-table has no row for o-17.

Correct answers: C and E

Explanation: A KTable represents changing keyed state, while a KStream represents individual records. When an update no longer satisfies KTable.filter, the resulting table receives a tombstone for that key, removing its previously matching value. Therefore, open-table no longer contains the order after the CLOSED update.

After conversion with toStream, KStream.filter simply discards nonmatching records. The CLOSED revision never reaches toTable, so that table receives no tombstone or replacement for the previously accepted OPEN revision. Its materialized state consequently remains stale at revision 1.

Filtering table updates preserves state-change semantics; filtering their stream representation can suppress the update needed to remove stale state.

Why each option fits or fails:

A. Storing revision 2 downstream ignores that the stream filter drops the CLOSED record before toTable receives it.

B. Treating the table filter as a stream filter overlooks the tombstone produced when the latest value fails the predicate.

C. The stream filter drops revision 2, so toTable receives no update that removes revision 1.

D. Removing the stream-derived row assumes a dropped record acts as a tombstone, but it produces no downstream record at all.

E. The nonmatching table update becomes a tombstone in the filtered table and removes the earlier row.


Question 12

Topic: Stateful stream processing and temporal patterns

An Apache Kafka Streams topology uses an ordinary KStream-KTable left join, with both inputs keyed by customer ID. Records are processed in this order:

  • The table records customer C7 as Silver.
  • Stream order O1 for C7 emits tier Silver.
  • The table changes C7 to Gold.

No further record for O1 arrives. Which interpretation follows from the join semantics?

Options:

  • A. O1 remains Silver; changing it requires explicit recomputation and a correction record.

  • B. O1 becomes Gold with a versioned table; as-of lookup re-emits prior results.

  • C. O1 becomes Gold; each table update re-evaluates and replaces earlier join results.

  • D. O1 becomes Gold at commit; caching merely delays the automatic join re-evaluation.

Best answer: A

Explanation: An ordinary stream-table join performs lookup-at-processing enrichment. Each stream record triggers a lookup against the table state available when that record is processed and produces a result. A later table update changes the state used for future stream records, but it does not retain, revisit, or re-emit earlier stream-side records. Therefore, the emitted result for O1 remains Silver.

Retroactive correction requires an explicit design, such as replaying or recomputing affected events under a defined temporal rule and emitting keyed corrections or replacements. A versioned table can provide bounded as-of lookup for newly processed stream records, but it still does not make later table updates automatically revise prior outputs.

Why each option fits or fails:

A. A table update changes future lookups but does not retrigger joins for previously processed stream records.

B. Versioned tables support bounded historical lookup, not automatic re-emission of prior results.

C. Table updates alter lookup state but do not trigger ordinary stream-table join outputs.

D. Cache and commit behavior affects output visibility, not which input side triggers evaluation.


Question 13

Topic: Stateful stream processing and temporal patterns

A Kafka Streams aggregation uses event timestamps and:

SessionWindows.ofInactivityGapAndGrace(
    Duration.ofMinutes(5),
    Duration.ofMinutes(3))

All acct-7 events use the same partition, and session state is retained for 30 minutes. Records arrive in this order on the same date:

ArrivalKeyEvent time
1acct-710:00
2acct-710:02
3acct-710:10
4acct-710:12
5acct-910:14
6acct-710:06

The fifth record advances task stream time to 10:14. What should the session aggregation do with the final record?

Options:

  • A. Accept 10:06 and merge all five events into [10:00, 10:12].

  • B. Accept 10:06 into [10:06, 10:12] and leave the earlier session unchanged.

  • C. Accept 10:06 into [10:00, 10:06] and leave the later session unchanged.

  • D. Reject 10:06 as late and leave the two existing sessions unchanged.

Best answer: A

Explanation: A session event can bridge multiple existing sessions when it falls within the inactivity gap of each. Here, 10:06 is four minutes after the earlier session’s 10:02 end and four minutes before the later session’s 10:10 start, so the candidate merged session spans 10:00 through 10:12.

Lateness is evaluated against that merged session’s end plus the five-minute inactivity gap and three-minute grace period. The session expires after stream time advances beyond 10:20. Because task stream time is only 10:14, the bridge is accepted and the two existing session entries are replaced by one merged aggregate.

Why each option fits or fails:

A. The bridge is within five minutes of both sessions, and the merged session remains open through stream time 10:20.

B. Merging only with the later session ignores that 10:06 is also within the inactivity gap of 10:02.

C. Merging only with the earlier session ignores that 10:06 is also within the inactivity gap of 10:10.

D. Rejecting based on 10:06 alone ignores that adjacent sessions extend the candidate session’s end to 10:12.


Question 14

Topic: Stateful stream processing and temporal patterns

A trade processor enriches raw trades with an instrument classification. Replays up to two years later must reproduce the classification known at the original decision time. Intraday changes and backdated corrections are possible.

Existing behavior: A compacted topic keyed by instrumentId retains the latest value, which is now R3.

Records for instrument X (UTC):

RevisionClassificationvalidFromrecordedAt
R1A09:0009:01
R2B10:0010:02
R3C10:0011:00

Trade T occurred at 10:05 and was originally decided at 10:06. R3 is a correction that was unknown at 10:06. Which design best supports reproducible enrichment?

Options:

  • A. Retain immutable daily reference snapshots for over two years; replay from the snapshot preceding each original decision.

  • B. Retain immutable reference revisions for over two years; replay with the valid-time cutoff and newest recorded revision.

  • C. Retain immutable reference revisions for over two years; replay with both valid-time and original recorded-time cutoffs.

  • D. Retain a compacted latest-reference topic for over two years; replay with the valid-time cutoff from each trade.

Best answer: C

Explanation: Reproducible enrichment requires both business valid time and system knowledge time. For trade T, R2 was effective at 10:05 and had been recorded by the 10:06 decision cutoff. R3 has the same effective time but was not recorded until 11:00, so using it would introduce future knowledge into the replay.

An immutable revision history must remain available for the full replay horizon. The lookup first restricts revisions to validFrom <= 10:05 and recordedAt <= 10:06, then selects the applicable revision. A latest-value compacted topic represents current keyed state and does not guarantee preservation of every historical revision.

Valid-time filtering alone cannot reproduce what the processor actually knew at an earlier cutoff.

Why each option fits or fails:

A. Daily snapshots fail because an intraday revision such as R2 may not appear in the preceding snapshot.

B. Newest recorded revision fails because it admits R3, which was unavailable at the original decision time.

C. The two cutoffs select R2, which was valid for the trade and known when the original decision occurred.

D. Compacted latest state fails because compaction can remove R2 even if the topic itself has a long retention period.


Question 15

Topic: Enterprise integration and business workflows

A Kafka normalization service consumes order events from web and partner systems and publishes one OrderChange contract.

Known conditions:

  • Event IDs are unique only within each source, and order IDs can overlap across sources.
  • Payloads contain source-only fields and distinct status vocabularies.
  • New status values can arrive before mapping rules are updated.
  • Operations must explain translations and replay inputs after rule updates.
  • Unmapped or invalid input must not appear as a valid business state.

Which TWO requirements should the design satisfy?

Options:

  • A. Attach source identity, schema and translator versions, and a durable original-payload reference to canonical records.

  • B. After successful mapping, discard source-only fields and the original payload, retaining no durable reference to that source evidence.

  • C. Store rejected inputs durably with the original payload, diagnostic reason, and replay reference on an error channel.

  • D. Translate every unfamiliar source status to UNKNOWN and publish it as a valid canonical state.

  • E. Deduplicate before translation by treating canonical order ID as the shared cross-source event identity.

Correct answers: A and C

Explanation: A Message Translator can produce a stable Canonical Data Model without destroying provenance. Because identifiers are source-scoped and source-only fields may later explain a translation, each canonical record needs source identity, relevant versions, and durable access to the original evidence.

Unmapped or invalid input is a distinct processing outcome, not an ordinary business state. A durable error channel should preserve the input and diagnostic context so operators can update mappings and replay it. Mapping an unfamiliar value to UNKNOWN would invent canonical meaning, while deduplicating by order ID would confuse entity identity with event identity and could merge unrelated sources.

Why each option fits or fails:

A. These fields preserve source-scoped identity and the evidence needed to explain each successful translation.

B. Discarding source-only fields and the original payload without a durable reference prevents later investigation of information omitted by the canonical contract.

C. A durable error outcome preserves failed normalization evidence and supports replay after mapping rules change.

D. Publishing unfamiliar statuses as UNKNOWN silently converts translation failures into apparently valid business events.

E. Order-level deduplication can suppress distinct events because order IDs overlap and do not identify source events.


Question 16

Topic: Enterprise integration and business workflows

An order service treats each committed status transition as accepted. It must preserve its database as the authority and emit one Kafka event for every accepted transition. Rolled-back transitions must emit nothing. Consumers can deduplicate repeated deliveries by stable event ID.

Observed flow:

BEGIN DB
UPDATE orders SET status = 'PAID' WHERE id = 742
COMMIT DB
producer.send(eventId = 'e-742-paid')

The process terminated after the database commit but before producer.send. XA transactions are unavailable, the database supports CDC, and multiple transitions may commit before a relay runs.

Which bounded design best closes this failure gap?

Options:

  • A. Write the order change and a per-order pending flag atomically, then publish the latest order state.

  • B. Publish the event in a Kafka transaction, then commit the order change after the broker acknowledges it.

  • C. Commit the order change, then use an idempotent Kafka producer with retries to publish the event directly.

  • D. Write the order change and an event-specific outbox row atomically, then publish the row through CDC.

Best answer: D

Explanation: A transactional outbox stores the business change and an event-specific publication record in the same database transaction. A rollback removes both; a commit preserves both. CDC then observes committed outbox rows and publishes them to Kafka. A relay failure can cause duplicate delivery, but the stable event ID allows consumers to deduplicate it.

This guarantee is not end-to-end atomicity between Kafka and the database. Instead, it removes the unrecorded interval between the database commit and the attempted send. Kafka producer idempotence and Kafka transactions cannot independently include the external database commit. Each transition needs its own outbox row because a single latest-state marker can collapse intermediate transitions.

Why each option fits or fails:

A. Pending latest state can collapse several committed transitions into one publication when the relay is delayed.

B. Kafka-first commit can expose an event even when the subsequent database transaction fails or rolls back.

C. Producer retries cannot recover a send that was never durably recorded before the process terminated.

D. The database transaction durably records each accepted transition and its publication intent, while CDC relays only committed outbox rows.


Question 17

Topic: Enterprise integration and business workflows

An order service currently commits each accepted command in one database transaction:

UPDATE Orders
SET status = 'CANCELLED', version = 5
WHERE order_id = 81 AND version = 4;

INSERT INTO Outbox(type, order_id, version, payload)
VALUES ('OrderStatusChanged', 81, 5, '{"status":"CANCELLED"}');

The outbox publishes to an order-keyed Kafka topic with cleanup.policy=delete and 14-day retention. Its stable integration payload intentionally omits internal decision details.

The service must retain accepted domain decisions for seven years, reconstruct any aggregate revision, and atomically reject commands with a stale expected version. Which design best satisfies these requirements?

Options:

  • A. Append domain events to a seven-year OrderEvents stream using an atomic expected-version check; derive the Orders projection and integration outbox.

  • B. Compact the integration topic by order and add aggregate versions; rebuild each aggregate from compacted values while retaining the database update and outbox flow.

  • C. Retain the integration topic for seven years, partition by order, and use event versions during replay; keep the database update and outbox flow.

  • D. Add seven-year temporal history to Orders and keep optimistic updates; rebuild aggregate states from row versions and retain the existing outbox flow.

Best answer: A

Explanation: Event sourcing requires an authoritative, append-only sequence of domain events with a defined retention contract. Each append must atomically confirm that the stream’s current version matches the command’s expected version. Projections such as Orders can then be rebuilt by replaying that sequence.

The existing outbox reliably transports integration records, but those records are not the domain authority. They omit internal decision details and currently expire after 14 days. Keeping the integration contract separate also allows external events to remain stable while the domain model evolves. The event append, projection update, and outbox insertion can share a database transaction where supported.

Publishing records to Kafka does not by itself make the application event sourced; authority, complete replay semantics, concurrency checks, and retention must all be explicit.

Why each option fits or fails:

A. The retained domain-event stream becomes authoritative, while atomic expected-version appends enforce concurrency before projections and integration records are derived.

B. Log compaction retains a keyed state view asynchronously and may remove intermediate transitions needed for revision replay.

C. Longer topic retention preserves transport records longer but neither restores omitted domain details nor atomically validates an append against the authoritative aggregate version.

D. Temporal row history records state versions, not the domain-event sequence required as the replay authority.


Question 18

Topic: Enterprise integration and business workflows

An organization is incrementally replacing a legacy customer application. Both applications store the same mutable customer fields, and both CDC connectors publish every database update.

Current behavior:

Legacy write -> legacy CDC e501 -> new DB apply
New DB apply -> new CDC e502 -> legacy DB apply -> ...

Constraints:

  • Customers must migrate in batches without an all-customer outage.
  • A routing registry stores each customer’s write owner and ownership epoch.
  • Applications can reject commands sent to a non-owner.
  • For a fenced customer, legacy prevents new commits and waits for previously accepted transactions to commit or abort; the CDC path can then expose a source checkpoint and its corresponding Kafka offsets.
  • Replication can preserve event identity and origin metadata.
  • Legacy screens still need current data after cutover.

Which design best supports the migration while preventing conflicting writes and replication feedback?

Options:

  • A. For each customer, switch routing and ownership immediately, drain pending legacy CDC afterward, then resolve overlap using each database’s update timestamp while filtering replication by provenance.

  • B. Keep legacy as owner for every customer during shadowing, wait until all CDC positions match, then switch routing globally and begin filtered reverse replication.

  • C. For each customer, keep both writers active during observation, preserve event IDs for deduplication, and resolve concurrent updates using the highest application version before selecting an owner.

  • D. For each customer, fence legacy writes, obtain a post-fence source checkpoint, drain legacy CDC through its corresponding Kafka offsets, switch routing and ownership, then copy new changes back while suppressing repeat replication by provenance.

Best answer: D

Explanation: A strangler migration should separate dual observation from write authority. While legacy owns a customer, the new application may consume legacy changes as a shadow. At cutover, fencing legacy commits, resolving previously accepted transactions, obtaining a source checkpoint, and draining CDC through its corresponding Kafka offsets ensure that no accepted legacy update remains behind the ownership boundary. The registry can then assign the new owner and epoch, allowing stale commands to be rejected.

Legacy may receive subsequent new-system changes as non-authoritative read copies. Replication must use preserved event identity and origin metadata to avoid forwarding those copies back, preventing a replication loop even though CDC publishes their database updates. Deduplication cannot reconcile two independently accepted business commands, and database timestamps do not reliably establish business causality. A global switch maintains one writer but fails the required batch migration.

Why each option fits or fails:

A. Switch before draining can apply an older legacy change after a new-owner write; database timestamps do not establish a safe cutover boundary.

B. Global ownership switch preserves single authority but does not meet the requirement to migrate customers in batches.

C. Dual write observation allows conflicting business commands; event-ID deduplication handles repeated events, not independent updates.

D. Resolving accepted writes and draining through the Kafka offsets corresponding to a post-fence source checkpoint establish a safe ownership boundary, while provenance filtering permits one-way read-copy replication without creating a feedback loop.


Question 19

Topic: Evolution security and data governance

Two Kafka services use distinct principals. Topics are pre-created, automatic topic creation is disabled, and all unlisted access is denied.

Required duties:

  • order-projector: consume orders.secure with group order-projector-prod, then produce to order.summary; no administration.
  • audit-exporter: consume order.summary with group audit-export-prod; no production or administration.

Current ACLs:

PrincipalOperationResource
order-projectorREADTopic orders.secure
order-projectorREADGroup order-projector-prod
order-projectorWRITE, CREATETopic order.summary
audit-exporterREADTopic order.summary
audit-exporterREADGroup order-projector-prod

Which TWO changes are required to satisfy the duties with least privilege?

Options:

  • A. Replace the audit exporter’s group READ with READ on audit-export-prod.

  • B. Grant the audit exporter READ on orders.secure because summaries originate there.

  • C. Grant the audit exporter WRITE on order.summary so it can commit offsets.

  • D. Remove the order projector’s CREATE grant on the pre-created order.summary topic.

  • E. Remove group READ from the order projector because topic READ authorizes consumption.

Correct answers: A and D

Explanation: Kafka authorizes ordinary consumer activity against both the topic and the consumer-group resource. The audit exporter has topic READ but lacks group READ for audit-export-prod; its grant on the projector’s group is both ineffective for its intended group and unnecessarily broad. The projector correctly has topic READ, group READ, and output-topic WRITE. Because order.summary is pre-created and automatic creation is disabled, its CREATE permission exceeds the projector’s stated duties.

Offset commits are consumer-group operations and do not require WRITE permission on the consumed topic. Least privilege therefore aligns each principal with its exact topics, operations, and consumer-group identity.

Why each option fits or fails:

A. An ordinary consumer needs group READ on the group identity it actually uses.

B. Reading a derived summary does not require permission to read its sensitive source topic.

C. Committing offsets uses authorization on the consumer group, not WRITE authorization on the topic.

D. Producing requires topic WRITE, while CREATE is unnecessary because the topic already exists.

E. Removing projector group READ would prevent normal group coordination and offset management even though topic READ remains.


Question 20

Topic: Evolution security and data governance

A team is validating a portable Kafka Streams design that uses a local state store and changelog instead of a remote database.

Architecture decision: Recovery must finish within 10 minutes after task assignment. The worst-case task has 18,000,000 changelog records, so the target environment must restore at least 30,000 records/s.

Migration rehearsal:

  • Steady-state latency and per-key ordering checks pass.
  • State decoders pass, and the restored state hash matches the baseline.
  • Task ownership remains stable with no retries.
  • Restoration sustains 12,000 records/s and takes 25 minutes.
  • Fetch and store-apply counters advance at similar rates.

Which ADR revision and next experiment best distinguish the remaining recovery bottleneck?

Options:

  • A. Record per-key ordering as failed; compare single-threaded and concurrent processing using the same keyed input.

  • B. Record assignment stability as failed; compare eager and cooperative rebalancing during another 10-minute recovery drill.

  • C. Record recovery throughput as failed; compare same-host no-op consumption with full restoration against 30,000 records/s.

  • D. Record state-format portability as failed; compare historical and current decoders using the same retained changelog.

Best answer: C

Explanation: The end-to-end recovery assumption is falsified in the target environment: restoring 18,000,000 records at 12,000 records/s requires 1,500 seconds, or 25 minutes, rather than the allowed 10 minutes. Stable ownership rules out rebalance delay for this run, while matching state hashes and successful decoding rule out an observed format failure.

Similar fetch and apply counters do not locate the bottleneck because state-store backpressure can limit fetching. Replaying the same changelog with identical host and fetch settings into a no-op consumer provides a baseline. A fast no-op replay implicates state application; equally slow results implicate the fetch, network, or broker path.

The ADR should retain the validated ordering and compatibility assumptions while marking recovery throughput as failed for this deployment.

Why each option fits or fails:

A. Ordering failure is unsupported because the keyed ordering check passed and does not explain the recovery duration.

B. Assignment instability is unsupported because task ownership remained stable and no retries occurred during restoration.

C. The recovery-rate assumption failed, and the paired replay distinguishes fetch-path limits from state-store application limits.

D. Format incompatibility is unsupported because decoding succeeded and the restored state matched the baseline hash.


Question 21

Topic: Testing observability and controlled operation

A team has this Kafka Streams topology:

builder.stream("raw", Consumed.with(Serdes.String(), readingSerde))
    .filter((key, value) -> value.quality() >= 80)
    .selectKey((key, value) -> value.deviceId())
    .mapValues(value -> (value.celsius() * 9) / 5 + 32)
    .to("accepted", Produced.with(Serdes.String(), Serdes.Integer()));

A converter-only unit test passed when a defect removed selectKey. The per-commit test must run without a broker and verify the complete deterministic topology using this single-partition fixture:

Input:  ingest-7, d1, 20 C, quality 95
Input:  ingest-8, d2, 15 C, quality 70
Input:  ingest-9, d1, 25 C, quality 90
Output: d1, 68
Output: d1, 77

Broker-based nightly tests already cover deployment behavior. Which test boundary and claim best meet the per-commit goal?

Options:

  • A. Recreate TopologyTestDriver between records; treat successful output as proof of rebalance and transaction recovery.

  • B. Run the fixture only in nightly broker tests; use consumer restarts to validate mapping and recovery.

  • C. Use TopologyTestDriver; assert the outputs and limit the claim to deterministic topology behavior.

  • D. Unit-test the converter with fixture values; infer that filtering and rekeying remain correctly wired.

Best answer: C

Explanation: TopologyTestDriver is the appropriate deterministic topology-test boundary. It can feed records through the built topology and inspect output keys, values, filtering, and order for the supplied sequence. This catches the removed selectKey defect that a converter-only unit test cannot detect.

The driver does not communicate with Kafka brokers or participate in consumer groups. It therefore does not prove network behavior, partition assignment, rebalances, transaction-coordinator behavior, or recovery from broker failures. Those properties require broker integration tests, while connector and external-system outcomes may require broader end-to-end tests.

The claim should remain limited to the topology behavior exercised by the fixture.

Why each option fits or fails:

A. Driver recreation does not execute Kafka’s rebalance protocol or transactional recovery mechanisms.

B. Nightly-only coverage can test broker behavior but leaves the required fast per-commit topology regression test absent.

C. The driver executes the complete topology against fixtures while making no claim about broker-mediated behavior.

D. Converter-only testing verifies arithmetic but bypasses the filter, key selection, and sink wiring where the reported defect occurred.


Question 22

Topic: Testing observability and controlled operation

A production Kafka route validates and enriches records from orders.in, then writes them to orders.ready. Separate production consumer groups reserve inventory through a database procedure and send confirmation email.

Operations must send a synthetic event through the deployed topics, processors, and production consumer logic without creating either external action. Transformations preserve defined payload fields but may discard Kafka headers. Which TWO requirements should the design adopt?

Ingress permits only an authorized operations principal to create the diagnostic marker and reserves a collision-free synthetic entity namespace; downstream consumers trust only that validated marker.

Options:

  • A. Wrap marked processing in Kafka transactions so database and email effects are rolled back on abort.

  • B. Run production consumers for marked events, replacing external adapter calls with auditable simulated outcomes.

  • C. Publish diagnostic events to a cloned topic chain, then compare its final records with production.

  • D. Route marked final records to a diagnostic consumer group instead of the production consumer groups.

  • E. Propagate a schema-defined diagnostic marker, correlation ID, and synthetic entity IDs through every transformation.

Correct answers: B and E

Explanation: A diagnostic event must remain distinguishable throughout the live route, and every side-effect boundary must honor that distinction. Because transformations may discard Kafka headers, the marker and correlation information belong in preserved payload fields. Synthetic entity IDs prevent the event from modifying state associated with a real order.

The deployed production consumers should still deserialize the event and execute their normal business rules. At the external adapter boundary, database and email operations are replaced with recorded simulated outcomes. This last-moment detour tests more of the real path while preventing business effects. A cloned route or separate consumer group bypasses production components, while Kafka transactions cannot roll back database procedures or email delivery.

Why each option fits or fails:

A. Kafka transaction coordinates Kafka records and offsets, not atomic rollback of database procedures or delivered email.

B. Intercepting marked events at the adapter boundary exercises production consumer logic without executing real business actions.

C. Cloned route tests a separate configuration and therefore cannot establish that the deployed live route works.

D. Separate consumer group bypasses the production consumers whose deserialization and business logic must be exercised.

E. Payload fields preserve diagnostic identity across the stated transformations, while synthetic IDs isolate the event from real entities.


Question 23

Topic: Architecture tradeoffs and integrated design cases

A retailer publishes OrderAccepted. Inventory, payment, and shipping services consume it in separate consumer groups and initiate work independently. Shipping has started before payment authorization, and failed payments have left inventory reserved.

Delivery conditions:

  • Outcome records are keyed by orderId and delivered at least once.
  • No transaction spans the services’ databases.
  • Services must remain asynchronous and independently deployable.
  • After a restart, processing must resume from the durable workflow state.
  • Operations require one authoritative status and an audit of transition decisions.

Scroll sideways if needed. Open full-size diagram in a new tab

Text description

Order acceptance starts inventory reservation. Payment authorization depends on a successful reservation. Shipment depends on authorization. Failed payment or no payment outcome within 10 minutes requires inventory release.

Which design best enforces this workflow?

Options:

  • A. Use an order-keyed process manager with durable state, reliable command dispatch, idempotent participants and outcome handling, and deadline ownership.

  • B. Chain predecessor events between services, with each service persisting its result and a scheduler releasing expired reservations.

  • C. Join reservation and payment events in Kafka Streams while services continue initiating their steps from OrderAccepted.

  • D. Publish compacted workflow-status records, allowing each service to claim and execute any currently permitted transition.

Best answer: A

Explanation: Kafka distributes records but does not itself own or enforce a multi-step business process. Each service consuming OrderAccepted independently can act before its prerequisites are satisfied. An order-keyed process manager maintains a durable state machine, records completed steps and deadlines, and emits the next command only when the required outcome has occurred. Idempotent outcome handling protects the state machine from redelivery. State transitions must be coupled reliably with command publication, such as through an applicable Kafka transaction or a transactional outbox, and stable command identities provide a basis for destination-side deduplication rather than protection by themselves.

The process manager can use Kafka-backed state, a database, or another durable store, but its ownership boundary must be clear. Kafka transactions alone cannot atomically include the independent service databases, so destinations must deduplicate commands or implement idempotent business operations. A stream join can correlate successes, but it does not govern services that already start directly from the original event.

Why each option fits or fails:

A. A durable process manager owns the workflow state, dependency decisions, retries, and timeout transition for each order. It couples persisted transitions with reliable command publication, while destinations deduplicate stable command IDs or make the commanded operation idempotent.

B. Distributed choreography can sequence events, but local records and a separate scheduler do not provide one authoritative workflow decision owner.

C. Success-event join may gate shipment, but payment still starts independently of the required reservation outcome.

D. Compacted status claims provide a state view, but compaction and producer idempotence do not serialize competing business transitions across services.


Question 24

Topic: Architecture tradeoffs and integrated design cases

A valuation service receives a EUR 100 trade at 16:30 UTC on June 30. The FX feed provides these versions, in USD per EUR:

VersionValid fromRecorded atRate
v1June 30, 16:00June 30, 16:021.08
v2June 30, 16:00July 2, 09:001.10

Version v2 is an accepted correction. A report generated July 1 at 02:00 includes facts valid by June 30 at 17:00 and known by its generation time. It must be reproducible by recalculation. The July 3 dashboard must use all accepted facts.

Which TWO conclusions are supported?

Options:

  • A. The service should overwrite v1 because v2 has the same business-valid start.

  • B. The service should value the report at USD 108 and the dashboard at USD 110.

  • C. The service should value both views at USD 110 because v2 is valid before the report cutoff.

  • D. The service should retain both versions with their actual valid and recorded times.

  • E. The service should value both views at USD 108 because the report knowledge cutoff is fixed.

Correct answers: B and D

Explanation: Business valid time identifies when a rate applies, while recorded time identifies when the service learned that version. The correction is valid for the trade’s time but was not known when the July 1 report was generated. An as-known query therefore applies both the business cutoff and the knowledge cutoff, selecting v1 and producing USD 108. The July 3 current view uses the latest accepted version valid at the trade time, producing USD 110. Retaining both versions with immutable recorded times supports both queries; overwriting or backdating the correction would introduce future-data leakage into the reproduced report.

Why each option fits or fails:

A. Overwrite v1 destroys the historical rate needed to recalculate what was known on July 1.

B. The report sees v1, while the later dashboard sees v2, producing EUR 100 times 1.08 and 1.10 respectively.

C. USD 110 for both ignores that the correction was recorded after the report’s knowledge cutoff.

D. Versioned temporal history preserves both current truth and the facts available at the report’s knowledge cutoff.

E. USD 108 for both incorrectly applies the historical report’s knowledge cutoff to the current dashboard.


Continue in the web app

Use IT Mastery for interactive Kafka Design Patterns practice with mixed sets, timed mocks, topic drills, explanations, and progress tracking.

Try Kafka Design Patterns on Web

Recall the key distinctions · Technical references · Report a question issue