Free CCDAK Practice Questions: Apache Kafka Developer

Try 24 original CCDAK practice questions with code, diagrams and explanations across Kafka clients, Streams, Connect, testing and observability.

These 24 original IT Mastery practice questions cover all six CCDAK areas. There are 18 single-answer questions and six Select TWO questions. Read the requested answer count, inspect the evidence, and try each question before revealing its explanation.

This short preview is designed for learning and does not predict a passing result. The sample length and interaction mix are editorial choices. The issuer’s printed topic percentages total 99%; practice normalizes their proportions and rounds to whole questions. These are not official Confluent exam questions, copied live-exam content or exam dumps.

For a correct guess, identify the fact that makes your selection defensible. Revisit unfamiliar concepts in the official documentation, then try fresh questions in the app; repeated attempts on this fixed page can reflect answer recognition.

Practice-set coverage

DomainOfficial rangeQuestions in this set
Apache Kafka Fundamentals23%5
Apache Kafka Application Development28%7
Apache Kafka Streams12%3
Kafka Connect15%4
Application Testing8%2
Application Observability13%3

Practice questions

Questions 1-24

Question 1

Topic: Application testing

A developer uses Apache Kafka 4.0 test libraries to test this topology:

KTable<String, String> history = builder
    .stream("events", Consumed.with(Serdes.String(), Serdes.String()))
    .groupByKey()
    .aggregate(
        () -> "",
        (key, value, aggregate) ->
            aggregate.isEmpty() ? value : aggregate + ">" + value,
        Materialized.<String, String, KeyValueStore<Bytes, byte[]>>as("history-store")
            .withKeySerde(Serdes.String())
            .withValueSerde(Serdes.String()));

history.toStream().to("updates",
    Produced.with(Serdes.String(), Serdes.String()));

A TopologyTestDriver receives these records one at a time on a one-partition input topic with increasing timestamps:

SequenceKeyValue
1acct-1open
2acct-1approved
3acct-2open
4acct-1Kafka null
5acct-1closed

The fourth value is a Kafka null value, not the text "null". After supplying all records, the test drains updates and queries history-store.

Which conclusion is INCORRECT?

Options:

  • A. Enabling the Streams cache and lengthening commit.interval.ms makes this driver emit only each key’s final update, matching production coalescing.

  • B. The store ends with acct-1=open>approved>closed and acct-2=open, with no state update caused by the null value.

  • C. Reversing the first two non-null acct-1 records changes both its intermediate updates and its final stored history.

  • D. The updates topic contains, in order, acct-1=open, acct-1=open>approved, acct-2=open, and acct-1=open>approved>closed.

Best answer: A

Explanation: A KGroupedStream aggregation skips records whose key or value is null, so the fourth input neither invokes the aggregator nor changes the state store. Each other record updates the corresponding table entry.

TopologyTestDriver processes and flushes each supplied input synchronously. Consequently, the output topic exposes every intermediate table update: two updates for the first account before the second account appears, followed by the final first-account update. The materialized store contains only the latest value for each key after all inputs have been processed.

The aggregation concatenates values, so record order affects both intermediate updates and the final stored string. Changing cache size or commit timing in a driver test does not simulate production cache coalescing; integration testing is needed to observe such runtime forwarding behavior.

Why each option fits or fails:

A. TopologyTestDriver processes and flushes each input synchronously; cache and commit-interval settings do not reproduce production update coalescing.

B. A grouped-stream aggregation skips the null-valued record, so the materialized store retains the latest aggregate from accepted records.

C. String concatenation is order-sensitive, so reversing those records produces approved>open>closed rather than open>approved>closed.

D. The driver flushes after each input, while the null value is skipped, producing the four listed updates in input order.


Question 2

Topic: Kafka Streams

A Kafka Streams 4.0 application uses application.id=order-counts-v1 and the following stateful topology:

orders.groupByKey().count(Materialized.as("count-store"));

Effective assumptions:

  • count-store is a persistent RocksDB store with changelog logging enabled.
  • num.standby.replicas=0.
  • Instance A owns the task for orders partition 2.
  • The task completed its last commit before the failure, with no later records processed.

Instance A fails, and its local disk becomes permanently unavailable. The Kafka brokers remain healthy. Instance B has an empty state directory and receives the task.

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

Text description

Orders partition 2 feeds a task on instance A. The task updates a local RocksDB count store and writes state changes to changelog partition 2 in Kafka. Instance A then fails with its disk unavailable, and the task is reassigned to instance B, which has empty local state.

Which sequence describes the normal Kafka Streams recovery path?

Options:

  • A. Recover the task by replaying orders through the committed input position to rebuild count-store, then processing subsequent input.

  • B. Recover the task by copying count-store from instance A’s RocksDB directory, then applying changelog records written after its checkpoint.

  • C. Recover the task by creating an empty count-store, then resuming orders at the committed input position because offsets encode prior counts.

  • D. Recover the task by rebuilding count-store from its changelog partition, then resuming orders at the committed input position.

Best answer: D

Explanation: A persistent RocksDB store is the task’s local working copy of state. With logging enabled, Kafka Streams also writes store changes to a partitioned internal changelog topic. Because instance A’s disk is lost and no standby replica exists, instance B creates a local store and restores it by consuming the corresponding changelog partition. Once restoration reaches the committed state, the task resumes ordinary input processing from its committed position.

Persistent local storage can make a restart faster when the same disk remains available, but it is not the only durable copy. Input offsets record positions rather than aggregate contents, and ordinary input-topic replay is not the normal restoration mechanism for this logged materialized store.

Why each option fits or fails:

A. With changelog logging enabled, Streams restores the aggregate from its changelog rather than rebuilding it through ordinary input-topic processing.

B. Instance A’s disk is permanently unavailable; reusable local state can accelerate recovery only when that storage remains accessible.

C. Committed offsets identify input positions, not aggregate values, so they cannot replace the lost contents of the state store.

D. The Kafka-backed changelog reconstructs the lost local store, while the committed input position determines where ordinary processing resumes.


Question 3

Topic: Application observability

A Java consumer group writes fulfillment rows from an orders topic. Each accepted order has a unique orderId, and exactly one durable fulfillment effect should exist per accepted order. The topic still retains the entire incident window.

Post-repair evidence:

  • Total group lag fell from 12,000 at 10:20 to 4,100 at 10:30 and 0 at 10:40. It remained 0 through 11:00.
  • From 10:40 through 11:00, consumption matched production, processing errors were 0, and no rebalances occurred.
  • Sample application traces showed resumed processing:
orderId=O-104 partition=1 offset=8821 result=APPLIED
orderId=O-220 partition=2 offset=3310 result=ALREADY_PRESENT
orderId=O-317 partition=2 offset=3344 result=APPLIED
  • The accepted-order ledger contains 12,000 distinct incident-window IDs.
  • The fulfillment database contains 12,001 incident-window rows but only 11,998 distinct orderId values.

Incident closure requires read-only evidence that every accepted order produced exactly one durable fulfillment effect. Which next verification most directly closes the remaining uncertainty?

Options:

  • A. Perform a three-way identity reconciliation across accepted order IDs, retained Kafka partition-offset records, and durable fulfillment rows, listing every zero-effect or multiple-effect ID.

  • B. Extend lag, consumption-rate, error-rate, and rebalance monitoring for another observation interval to confirm that the application remains operationally stable.

  • C. Compare accepted-order and fulfillment-row totals in five-minute buckets throughout the incident and recovery periods.

  • D. Compare each partition’s committed group offset with its log end offset and inspect consumer logs around the final committed positions.

Best answer: A

Explanation: Declining lag followed by sustained zero lag shows that the consumer group caught up to the topic’s current log ends. Normal rates and sample traces also show that traffic resumed. These operational signals do not establish business recovery: committed progress and successful processing logs are not durable proof that every external database effect exists exactly once.

The aggregate reconciliation already reveals a mismatch, but totals cannot identify its cause because missing and duplicate effects can offset one another. A read-only, identity-level reconciliation must map each accepted orderId to its retained Kafka partition and offset, then count durable fulfillment rows for that identity. This distinguishes accepted orders never published, Kafka records with no durable effect, and records producing multiple effects, providing the evidence needed for incident closure.

Why each option fits or fails:

A. This links each business expectation to its Kafka record and durable effect cardinality, directly identifying unpublished, missing, and duplicate outcomes.

B. Longer operational stability confirms resumed traffic but cannot determine which accepted IDs have missing or duplicate durable effects.

C. Aggregate buckets can conceal offsetting discrepancies when one order is missing and another has multiple fulfillment rows.

D. Offset positions can confirm catch-up through Kafka, but they do not prove that each corresponding database effect exists exactly once.


Question 4

Topic: Application development

A team promotes one Java producer image through development and production.

  • Development injects development SASL credentials and mounts a truststore containing the development cluster’s private CA.
  • Production injects production SASL credentials and mounts a truststore containing the production cluster’s different private CA.
  • Credentials must not be stored in the image, source repository, or container filesystem.
  • Neither private CA is present in the JVM default truststore.
  • Broker certificates contain the correct DNS names.

The application loads this file directly with Properties.load:

security.protocol=SASL_SSL
sasl.mechanism=PLAIN
sasl.jaas.config=org.apache.kafka.common.security.plain.PlainLoginModule required username="${KAFKA_USERNAME}" password="${KAFKA_PASSWORD}";
ssl.truststore.location=${KAFKA_TRUSTSTORE_PATH}
ssl.truststore.password=${KAFKA_TRUSTSTORE_PASSWORD}
ssl.endpoint.identification.algorithm=https

Production startup fails with NoSuchFileException: ${KAFKA_TRUSTSTORE_PATH}. Which deployment change fixes the failure while satisfying the security requirements?

Options:

  • A. Render the injected values into a generated client.properties file, load that file, reference the environment-specific truststore, and retain HTTPS endpoint identification.

  • B. Read the injected credentials in Java, remove the truststore properties to use the JVM default truststore, and retain HTTPS endpoint identification.

  • C. Read the injected values in Java, set the Kafka properties at runtime, reference the environment-specific read-only truststore, and retain HTTPS endpoint identification.

  • D. Read the injected credentials and truststore settings in Java, reference the environment-specific truststore, and disable HTTPS endpoint identification.

Best answer: C

Explanation: Properties.load preserves ${...} as literal text; it does not automatically substitute environment variables. The application must therefore read the injected values and place them into the Kafka client configuration before constructing the producer.

SASL credentials authenticate the client, while TLS trust and endpoint identification authenticate the broker. Each deployment should reference its own mounted truststore containing the appropriate private CA. Keeping ssl.endpoint.identification.algorithm=https also verifies that the broker’s DNS name matches its certificate. Constructing the configuration in memory avoids baking credentials into the shared image or writing them to a generated properties file.

Why each option fits or fails:

A. Rendering resolves the placeholders but writes the SASL credentials to the container filesystem, violating the stated secret-handling requirement.

B. The default JVM truststore does not contain the stated private production CA, so TLS certificate-chain validation still fails.

C. Java resolves the deployment values before creating the producer, while the mounted CA truststore and hostname verification authenticate the production broker.

D. The truststore validates the issuing CA, but disabling endpoint identification removes verification that the certificate belongs to the broker’s DNS name.


Question 5

Topic: Kafka Connect

A Kafka Connect HTTP sink task has this state:

  • Assigned partition: payments-2
  • Last recorded next offset: 730
  • Records delivered in the current batch: offsets 730 through 734
  • Failure: All five HTTP writes returned success, but the task terminated before Connect recorded offset 735.
  • Each Kafka record represents a distinct destination row. Different records may have the same paymentId or identical payloads.
  • The current sink generates a new HTTP request identifier for every attempt.

Destination contract: The first request using an idempotency key creates one row. A retry using the same key and byte-identical body returns success without creating another row. Reusing the key with a different body is rejected.

The five rows from this failed test will be reconciled separately. The team will then repeat the test from a clean destination using a corrected connector. The design must produce exactly one destination row per Kafka record and recover automatically without loss. Which TWO actions must the corrected design use from the first delivery and during recovery?

Options:

  • A. Derive each HTTP idempotency key from a hash of the request body, allowing identical payloads to share one destination identity.

  • B. From the first delivery, derive each HTTP idempotency key from the Kafka topic, partition, and offset, and reuse the original request body.

  • C. After the same failure in the corrected test, resume from offset 730, resend records with their original idempotency keys, and treat duplicate-success responses as successful delivery.

  • D. Enable idempotence for the connector’s Kafka producer and generate a new HTTP request identifier whenever a record is retried.

  • E. Derive each HTTP idempotency key from paymentId, allowing all records for the same payment to share one destination identity.

Correct answers: B and C

Explanation: Kafka Connect’s consumed-offset progress is separate from successful external writes. The committed value 730 means offset 730 is the next record to consume. Because the task failed before recording 735, offsets 730 through 734 can be replayed after restart even though their HTTP writes succeeded.

In the corrected test, safe recovery requires a stable destination identity for each Kafka record from its first delivery. Assigning new keys only after the original failed test would not match the keys used for its five successful writes, so those existing rows require separate reconciliation. The tuple of topic, partition, and offset uniquely identifies the record within Kafka and remains unchanged during replay. Resending the same body with that identity allows the endpoint to return success without inserting another row. Connect can then record progress through offset 735.

A business identifier or body hash is unsuitable because separate Kafka records can share those values. Kafka producer idempotence also does not extend to an HTTP destination.

Why each option fits or fails:

A. Distinct Kafka records may have identical payloads, so a body hash could incorrectly deduplicate records that each require a row.

B. Topic, partition, and offset provide a stable identity for one Kafka record, allowing the destination to suppress duplicates safely.

C. Offset 735 was not recorded, so replay is expected; accepting idempotent duplicate responses lets Connect safely advance progress afterward.

D. Kafka producer idempotence does not deduplicate external HTTP writes, and new request identifiers would permit duplicate destination rows.

E. Multiple distinct records may share a payment identifier, so this identity could incorrectly collapse valid destination rows.


Question 6

Topic: Kafka Connect

A sink connector exports Kafka records to an external service.

Baseline contract:

  • Export only orders.v1.
  • Use string keys and schemaless JSON values.
  • Run no more than three tasks.
name=order-json-export
connector.class=com.example.ExportSinkConnector
tasks.max=3
topics=orders.v1
key.converter=org.apache.kafka.connect.storage.StringConverter
value.converter=org.apache.kafka.connect.json.JsonConverter
value.converter.schemas.enable=false

The connector supports these topic-selection properties:

  • topics selects an explicit comma-separated list.
  • topics.regex uses a Java regular expression against full topic names and periodically discovers new matches.
  • Exactly one of topics and topics.regex may be configured.

Existing topics include orders.v1, orders.v2, ordersXv9, and inventory.v1.

Changed condition: The connector must now export every current and future topic whose name is orders.v followed by one or more digits. All other contract requirements remain unchanged.

Which configuration revision satisfies the changed contract?

Options:

  • A. Replace topics with topics=orders.v1,orders.v2.

  • B. Remove topics and set topics.regex=orders.v[0-9]+.

  • C. Remove topics and set topics.regex=orders[.]v[0-9]+.

  • D. Keep topics=orders.v1 and add topics.regex=orders[.]v[0-9]+.

Best answer: C

Explanation: Dynamic topic discovery requires topics.regex rather than an explicit topics list. The expression orders[.]v[0-9]+ matches the required namespace: [.] represents a literal period, and [0-9]+ represents one or more digits. It therefore includes orders.v1, orders.v2, and future names such as orders.v10, while excluding ordersXv9.

Because topics and topics.regex are mutually exclusive, the explicit topics property must be removed. The connector identity, task ceiling, and converter settings remain unchanged because the revised contract changes only topic selection. tasks.max=3 remains a ceiling, while the existing converters continue producing string keys and schemaless JSON values.

Why each option fits or fails:

A. An explicit list selects the current topics but will not automatically include future numbered order topics.

B. The unescaped period matches any character, so the expression also selects the unrelated ordersXv9 topic.

C. The character class matches the literal period, while the digit expression includes current and future numbered order topics.

D. The connector permits only one topic-selection property, so configuring both selection modes is invalid.


Question 7

Topic: Application development

A team deploys an Apache Kafka 4.0 Java application with three logical processing slots. Each slot consumes input records, writes derived records transactionally, and sends the consumed offsets into the same Kafka transaction.

Deployment facts:

  • ordinal identifies a logical slot and remains stable when that slot restarts.
  • podUid identifies one process incarnation and changes after every restart.
  • All slots must divide one input subscription.
  • A replacement process must fence any lingering producer from the same logical slot.
  • Client metrics must distinguish both the process incarnation and client role.

The application uses this configuration template:

consumerProps.put(ConsumerConfig.CLIENT_ID_CONFIG,
                  clientBase + "-consumer");
consumerProps.put(ConsumerConfig.GROUP_ID_CONFIG, groupId);
producerProps.put(ProducerConfig.CLIENT_ID_CONFIG,
                  clientBase + "-producer");
producerProps.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG,
                  transactionalId);

Which identifier mapping satisfies these requirements?

Options:

  • A. clientBase = "invoice-" + podUid; groupId = "invoice-" + ordinal; transactionalId = "invoice-tx-" + ordinal

  • B. clientBase = "invoice-" + podUid; groupId = "invoice-v1"; transactionalId = "invoice-tx-" + ordinal

  • C. clientBase = "invoice-" + podUid; groupId = "invoice-v1"; transactionalId = "invoice-tx-" + podUid

  • D. clientBase = "invoice-" + podUid; groupId = "invoice-v1"; transactionalId = "invoice-tx-v1"

Best answer: B

Explanation: Each identifier serves a different lifecycle boundary. client.id labels client activity for logs and metrics, so incorporating the process-specific podUid distinguishes incarnations; the role suffix distinguishes the consumer from the producer.

All consumers must share one group.id to divide the topic partitions within a single consumer group. Giving each slot another group ID would make each slot consume independently.

A transactional.id must be unique among concurrently valid transactional producers but stable when the same logical producer restarts. The stable ordinal provides that identity. When a replacement initializes transactions with the same transactional ID, Kafka advances the producer epoch and fences a lingering predecessor. A pod UID changes on restart and therefore cannot provide this fencing continuity.

Why each option fits or fails:

A. Using a distinct group ID for each slot makes every slot independently consume the subscription rather than sharing its partitions.

B. The process-specific client base supports observability, the shared group divides partitions, and the stable per-slot transactional ID enables fencing across restarts.

C. Changing the transactional ID on every restart prevents the replacement from fencing a lingering producer that used the previous process identity.

D. Sharing one transactional ID across active slots causes their producers to fence one another instead of operating as distinct transactional producers.


Question 8

Topic: Kafka fundamentals

A Java consumer processes account updates. Updates for the same account must reach an external ledger in ascending sequence order, while different accounts may be processed concurrently.

Broker evidence: The topic has three partitions, uses account IDs as keys, and its partition count remained stable.

partition  offset  key     sequence
0          40      acct-A  8
0          41      acct-A  9
1          12      acct-B  3

Application evidence: One consumer owns partitions 0 and 1. It submits records in poll iteration order to a three-thread executor. Each task writes to the ledger when it completes.

elapsed  event
0 ms     dispatch partition 0, offset 40, acct-A sequence 8
1 ms     dispatch partition 0, offset 41, acct-A sequence 9
2 ms     dispatch partition 1, offset 12, acct-B sequence 3
40 ms    ledger write completes for acct-B sequence 3
120 ms   ledger write completes for acct-A sequence 9
280 ms   ledger write completes for acct-A sequence 8

Which TWO statements are supported when the evidence sources are considered together?

Options:

  • A. Enabling idempotence on the upstream producer will make the ledger writes complete in partition-offset order.

  • B. The application should serialize tasks within each assigned partition while running different partition queues concurrently.

  • C. Increasing the topic’s partition count will let acct-A use more workers while preserving its sequence order.

  • D. Kafka retained acct-A partition order, while concurrent execution reversed the account’s ledger side-effect order.

  • E. The application can preserve acct-A ordering by committing offsets only after every task from the poll completes.

Correct answers: B and D

Explanation: Kafka guarantees record order only within a partition. Partition 0 contains acct-A sequence 8 at offset 40 before sequence 9 at offset 41, and the application dispatched them in that order. Kafka does not require concurrently executed tasks to finish in dispatch order, however. Sequence 9 therefore reached the external ledger before sequence 8.

Serializing work by assigned partition preserves partition order while still allowing partitions 0 and 1 to run concurrently. Committing after all tasks finish can improve processing recovery semantics, but it cannot reorder side effects that already occurred. Producer idempotence concerns duplicate records caused by producer retries, not downstream execution order. Adding partitions also does not provide ordered parallelism for one key and can change that key’s future partition mapping.

Why each option fits or fails:

A. Producer idempotence suppresses qualifying retry duplicates; it does not control consumer worker scheduling or external side-effect completion.

B. A serial queue for partition 0 preserves its offset order, while partition 1 can execute concurrently without violating the per-account requirement.

C. Same-key records do not gain ordered parallelism across partitions, and changing the partition count can remap future records for that key.

D. Offsets 40 and 41 establish the partition order, but sequence 9 reached the ledger before sequence 8 because its task completed first.

E. Delayed commits affect restart position but do not change the observed completion order of tasks within the polled batch.


Question 9

Topic: Application development

A one-time replay tool uses the Apache Kafka 4.0 Java consumer.

Consumer setup:

group.id=orders-service
enable.auto.commit=false
auto.offset.reset=latest

The tool calls subscribe() and then poll(Duration) until assignment completes. It discards that initial batch without processing it.

Observation after assignment:

Target timestamp: 10:00
Assignment: orders-0, orders-1
orders-0: beginning=100, committed=150, offsetForTimes=120, capturedEnd=180
orders-1: beginning=50,  committed=80,  offsetForTimes=null, capturedEnd=90

The tool must replay from the target timestamp through the captured end offsets, excluding records appended later. Afterward, the normal service must resume from the existing committed offsets.

The normal service is stopped for the replay, with no other members or commit writers in this group. Assignment, captured bounds and retention remain stable. Kafka record timestamps are nondecreasing with offset in each partition, and the time lookup uses those timestamps.

Which TWO actions correctly implement this bounded replay?

Options:

  • A. Process only offsets below 180 and 90, then close without committing the replay positions.

  • B. After assignment, seek orders-0 to 120 and orders-1 to 90 before the next replay poll.

  • C. Seek orders-0 to 120 and orders-1 to 90 immediately after subscribe(), before polling for assignment.

  • D. On completion, commit offsets 180 and 90 before closing the replay consumer.

  • E. After assignment, seek orders-0 to 120 and leave orders-1 at its committed offset 80.

Correct answers: A and B

Explanation: A Java consumer may seek only partitions in its current assignment. A seek changes the in-memory position used by the next fetch; it does not change the committed group offset. For orders-0, offsetsForTimes identifies 120 as the replay start. For orders-1, the null result means no available record has a timestamp at or after the target, so seeking to the captured end of 90 produces an empty bounded range.

End offsets are exclusive, so the tool processes orders-0 offsets 120 through 179 and no orders-1 records. It must ignore records at or beyond the captured ends if later polls include newly appended data. Because auto-commit is disabled and the replay consumer makes no manual commit, closing it leaves commits 150 and 80 unchanged. A later normal consumer in the same group resumes from those committed positions; auto.offset.reset=latest is irrelevant while valid commits exist.

Why each option fits or fails:

A. The captured end offsets bound the replay, while omitting commits preserves offsets 150 and 80 for the normal service’s restart.

B. Seeking assigned partitions changes their next fetch positions, and the captured end is appropriate where no offset exists at or after the target timestamp.

C. A subscribed consumer cannot seek these partitions until group assignment has made them part of its current assignment.

D. Committing the captured ends would replace the existing group progress and make the normal service resume from 180 and 90.

E. The valid commit at 80 prevents offset reset from applying, so this would replay pre-target records from orders-1.


Question 10

Topic: Application development

A Java application uses the Apache Kafka 4.0 producer.

Application and effective settings:

  • acks=all
  • enable.idempotence=true
  • Serialization succeeds, topic metadata is cached, and the producer buffer has capacity.
ProducerRecord<String, String> record =
    new ProducerRecord<>("orders", 2, "order-17", "created");

Future<RecordMetadata> result = producer.send(record, (metadata, error) -> {
    if (error == null) {
        log("callback-success " + metadata.topic() + "-" +
            metadata.partition() + "@" + metadata.offset());
    } else {
        log("callback-failure " + error.getClass().getSimpleName());
    }
});

log("send-return");
log("future-immediate=" + result.isDone());
RecordMetadata confirmed = result.get();
log("future-get " + confirmed.topic() + "-" +
    confirmed.partition() + "@" + confirmed.offset());

Observed trace:

12:00:00.100 [main] send-return
12:00:00.100 [main] future-immediate=false
12:00:00.118 [kafka-producer-network-thread] callback-success orders-2@417
12:00:00.119 [main] future-get orders-2@417

Which TWO conclusions are supported by the application and trace?

Options:

  • A. The successful callback and completed future report the acknowledged result under acks=all, including the assigned partition and offset.

  • B. The false isDone result means the record had not yet entered the producer accumulator and could not be transmitted.

  • C. Returning from send establishes that the leader and required in-sync replicas have already acknowledged offset 417.

  • D. Calling get initiated network transmission, while the earlier callback success was only a provisional notification.

  • E. Returning from send indicates acceptance for asynchronous transmission, but it does not yet confirm broker acknowledgment.

Correct answers: A and E

Explanation: With normal serialization, cached metadata, and available buffer space, send performs synchronous preparation and appends the record to the producer accumulator. It then returns a Future while the producer’s network thread handles batching, transmission, acknowledgments, and any supported retries.

Therefore, returning from send confirms neither broker receipt nor durability. The immediate isDone=false observation demonstrates that publication was still pending. The later successful callback and Future.get() result represent the same successful terminal outcome and provide the resulting partition and offset. With acks=all, that outcome follows the acknowledgment required from the leader and in-sync replicas. Calling get only blocks the application thread until completion; it does not start the send.

Why each option fits or fails:

A. Both completion mechanisms report successful metadata only after the producer receives the acknowledgment required by the effective acks setting.

B. An incomplete future means publication is unfinished, not that the successfully returned send failed to enqueue the record.

C. The future was still incomplete when send returned, so the required broker acknowledgment had not yet become application-visible.

D. The producer transmits asynchronously after send; get waits for the terminal result rather than initiating transmission.

E. The incomplete future immediately after send returned shows that local acceptance and broker-confirmed completion were separate events.


Question 11

Topic: Kafka fundamentals

A team uses Avro with Schema Registry compatibility set to BACKWARD. Version 1 is the latest registered schema.

Version 1:

{
  "type": "record", "name": "OrderTotal",
  "fields": [
    {"name": "orderId", "type": "string"},
    {"name": "currency", "type": "string"},
    {"name": "totalCents", "type": "long"}
  ]
}

Proposed version 2:

{
  "type": "record", "name": "OrderTotal",
  "fields": [
    {"name": "orderId", "type": "string"},
    {"name": "totalCents", "type": "double"},
    {"name": "sourceSystem", "type": "string", "default": "UNKNOWN"}
  ]
}

Schema Registry accepts version 2. The upgraded replay application uses version 2 as its reader schema and must process historical version 1 records. A fraud application remains on version 1 and must consume new version 2 records.

The business contract requires an integer number of minor currency units and an explicit currency. The new producer instead plans to write major-unit values such as 19.99 into totalCents.

Which TWO findings show why successful BACKWARD registration is insufficient to approve the rollout?

Options:

  • A. The version 2 replay application cannot resolve version 1 records because historical records lack sourceSystem.

  • B. The version 1 fraud application cannot resolve version 2 records using its existing reader schema.

  • C. The proposed event no longer preserves the amount and currency meaning required by downstream processing.

  • D. The BACKWARD result confirms that both version 1 and version 2 readers can consume records from either schema version.

  • E. The version 2 replay application cannot resolve the historical long value as a double.

Correct answers: B and C

Explanation: Avro BACKWARD compatibility asks whether the new schema can read data written with the preceding schema. Here, the version 2 reader can promote long to double, ignore the old currency field, and use the sourceSystem default. Registration can therefore succeed.

That result does not establish forward compatibility. The version 1 reader expects a required currency field that version 2 no longer writes, and Avro cannot narrow a writer’s double to a reader’s long.

Schema compatibility also addresses structural resolution rather than business meaning. Writing 19.99 as major currency units changes the meaning of a field named totalCents, while removing currency violates the stated downstream requirement. A rollout review must therefore assess reader direction and semantic invariants in addition to registry acceptance.

Why each option fits or fails:

A. The version 2 reader supplies UNKNOWN from the declared default when a version 1 writer record lacks sourceSystem.

B. The version 2 writer omits required currency and writes double, neither of which can be resolved by the version 1 reader.

C. A fractional major-unit value contradicts the minor-unit meaning of totalCents, and removing currency eliminates required interpretation context.

D. BACKWARD evaluates the new reader against prior writer data; it does not guarantee that an old reader can process new writer data.

E. Avro permits promotion from a writer’s long to a reader’s double, so this change does not prevent version 1 replay.


Question 12

Topic: Kafka fundamentals

A team increases a topic from 4 partitions to 8 while its Java producer and consumer group remain active. No client configuration or code changes are made.

Application assumptions:

  • The producer’s custom partitioner uses the serialized customer key and current partition count.
  • The consumer maintains customer state and assumes same-key records are processed in send order.
  • The group has two active workers.

Observed sequence:

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

Text description

The producer appends E100 with key C7 to partition 1. The topic owner then increases the topic from four to eight partitions. The coordinator assigns partition 1 to Worker A and partition 5 to Worker B, resolving partition 5’s initial position. E100 remains pending because Worker A has a backlog. The producer then appends E101 with the same key to partition 5, where it is available to Worker B.

E100 was sent before E101. The initial position for partition 5 was resolved before E101 was appended. Which outcome is supported by Kafka semantics?

Options:

  • A. Kafka moves E100 from partition 1 to partition 5 before making E101 consumable, preserving customer order.

  • B. Worker B may process E101 before Worker A processes E100 because the records are now in different partitions.

  • C. The coordinator assigns partitions 1 and 5 to the same worker because both contain records with customer key C7.

  • D. The partition increase invalidates the committed position for partition 1, causing offset reset before E100 can be read.

Best answer: B

Explanation: Kafka ordering applies within an individual partition. Increasing a topic’s partition count adds empty partitions but does not move existing records or rewrite their offsets. A partitioner that incorporates the current partition count may therefore route a key differently after the increase.

Here, E100 remains in partition 1 while the later E101 is appended to partition 5. The consumer-group coordinator assigns partitions rather than customer keys, so the two partitions can have different owners. Because Worker A has a backlog, Worker B can process E101 first even though the producer sent it later. Existing committed positions for partition 1 remain valid; only newly added partitions require an initial consumer position. The partition increase is operationally permitted but breaks the application’s assumption that same-key processing order will continue across the change.

Why each option fits or fails:

A. Increasing the partition count does not redistribute records already stored in existing partitions.

B. Kafka preserves order within a partition, but it does not impose processing order across partitions owned by separate workers.

C. Consumer-group assignments operate on partitions and do not co-locate partitions based on their record keys.

D. Adding partitions does not invalidate committed offsets for partitions that already existed.


Question 13

Topic: Kafka fundamentals

A Java consumer builds a local materialized view from keyed records in a single-partition topic.

  • At shutdown, the view contained K and reflected records through offset 120, and the committed next offset was 121.
  • While the consumer was offline, key K received a tombstone at offset 240.
  • The topic uses cleanup.policy=compact without time-based deletion.
  • delete.retention.ms is 24 hours, and the consumer was offline for 10 days.
  • Compaction has removed the tombstone and all older records for K.
  • The log-start offset remains 0, but the surviving application records have sparse offsets from 300 through 900; the end offset is 901. These records include the latest value of every live key.
  • No records are appended and no further cleanup occurs during the bounded rebuild.

Which recovery action is sufficient to produce a complete current materialized view?

Options:

  • A. Discard the persisted view, seek to offset 901, and rebuild state using only subsequently appended records.

  • B. Keep the persisted view, reset to offset 901, and process only records appended after recovery.

  • C. Keep the persisted view, reset to offset 300, and apply the retained records to the existing state.

  • D. Discard the persisted view, reset to offset 300, and replay the retained log into an empty state.

Best answer: D

Explanation: Log compaction retains the latest value for each live key, but deletion tombstones need only remain available for the configured tombstone retention period. A consumer that keeps an old local view must observe later tombstones to remove deleted keys. Here, the persisted view contains K, but the tombstone that deleted it has already been removed.

Replaying from offset 300 on top of that stale view would therefore leave K incorrectly present. Rebuilding from an empty state avoids this problem: live keys are reconstructed from their latest retained records, while deleted keys remain absent. The committed offset identifies saved consumer progress, but it does not guarantee that the corresponding log position or every required deletion event is still available.

Compaction removes records without renumbering offsets. Because the log-start offset is still 0, seeking to 121 would be valid and would advance to the next surviving record; it must not be rejected merely because the first surviving record is at 300. The incorrect empty-state plan instead starts at 901 and misses all retained live values.

Why each option fits or fails:

A. Starting empty at the end omits live keys whose latest values are already in the retained log and may never be updated again.

B. Starting at the end preserves stale K and also skips every retained update produced during the downtime.

C. The existing state still contains K, and replay cannot remove it because its tombstone is no longer retained.

D. Starting empty omits deleted keys while the compacted log supplies the latest retained value for each currently live key.


Question 14

Topic: Application development

A Java application writes Avro values to the payments topic. The same partition contains records produced with both schema versions.

Schema Registry setup:

  • Subject: payments-value, using TopicNameStrategy
  • Version 1, schema ID 17: fields transaction_id as string and amount_cents as long
  • Version 2, schema ID 42:
{
  "type": "record",
  "name": "Payment",
  "namespace": "com.example",
  "fields": [
    {"name": "transaction_id", "type": "string"},
    {"name": "amount_cents", "type": "long"},
    {"name": "memo", "type": ["null", "string"], "default": null}
  ]
}
  • Compatibility mode: BACKWARD

Producer configuration:

  • value.serializer=io.confluent.kafka.serializers.KafkaAvroSerializer
  • auto.register.schemas=false
  • use.latest.version=false
  • The producer supplies a GenericRecord using the registered version 2 schema.

Consumer configuration:

  • value.deserializer=io.confluent.kafka.serializers.KafkaAvroDeserializer
  • specific.avro.reader=false
  • The consumer uses the same Schema Registry.

Which conclusion is INCORRECT?

Options:

  • A. The producer resolves version 2 to schema ID 42 and writes Confluent framing followed by the Avro-encoded value.

  • B. The consumer uses each record’s embedded schema ID to resolve its writer schema and returns a GenericRecord.

  • C. The consumer can ignore the embedded schema ID and decode every record using version 2 because BACKWARD establishes one topic-wide payload layout.

  • D. The broker stores the serialized values as bytes and does not enforce the registry’s compatibility mode when records are produced.

Best answer: C

Explanation: A Confluent-framed Avro value contains a magic byte, a four-byte schema ID, and the Avro payload. The full schema is retained in Schema Registry rather than copied into every record. With auto.register.schemas=false and use.latest.version=false, the producer looks up its exact version 2 schema and uses ID 42.

The broker remains unaware of Avro and Schema Registry compatibility rules. During deserialization, the client uses each record’s ID to resolve the writer schema, potentially from its local cache. This matters because the partition may contain records written with IDs 17 and 42.

BACKWARD compatibility means the version 2 schema can read data written with version 1 when Avro reader-writer resolution is performed. It neither rewrites old payloads nor makes the latest schema sufficient to identify every record’s original encoding.

Why each option fits or fails:

A. With automatic registration and latest-version selection disabled, the serializer looks up the exact registered schema and uses its existing ID.

B. The deserializer reads the framed ID to resolve the appropriate writer schema, while specific.avro.reader=false selects generic Avro records.

C. Backward compatibility supports reader-writer schema resolution but does not make different writer encodings identical or eliminate the embedded schema ID.

D. Kafka brokers treat record values as bytes; schema compatibility is handled by Schema Registry and client serializers rather than the broker.


Question 15

Topic: Application observability

A Kafka Streams application runs groupByKey().count() over a four-partition input topic, then writes the resulting KTable updates to a compacted output topic. KTable caching is enabled, and the topology is unchanged between tests.

Workload and task evidence:

  • The same sustained workload is used in both tests.
  • One key represents 90% of input records and always routes to partition 2.
  • Four active tasks exist in both tests.
  • After the change, the hot task is 98% busy, other tasks are at most 30% busy, and two threads have no active task.
  • No rebalances or state restorations occur during measurement.
Metric over 10 minutes2 threads6 threads
Input throughput18,200 records/s19,100 records/s
Input-to-forward p99 latency210 ms390 ms
State-store puts18,100 operations/s19,000 operations/s
RocksDB write-stall time0 ms/s0 ms/s
Output throughput6,800 records/s4,900 records/s

After each test drains, the latest output value for every key matches the expected count.

Which conclusion and next change are best supported by this evidence?

Options:

  • A. The hot key is saturating one task; additional threads cannot parallelize it. Treat the final counts as correct and consider a semantics-preserving shard-and-combine aggregation.

  • B. The four-task ceiling is the main limit; increase the input topic to six partitions while retaining the existing keys, giving every stream thread an active task.

  • C. Output-topic capacity is the main limit; increase its partition count and disable cache coalescing so that output throughput matches input throughput.

  • D. State-store write saturation is the main limit; enlarge the RocksDB cache before adding threads, and regard the reduced output rate as incomplete aggregation.

Best answer: A

Explanation: Kafka Streams parallelism is based on tasks, which are derived from input partitions. A record key is processed by one partition task, so adding threads cannot divide the work for a single hot key. Here, the hot task is nearly saturated while other tasks and two threads have spare capacity, explaining the small throughput gain and increased latency.

The state-store evidence does not indicate a RocksDB bottleneck: put activity follows processing throughput, and no write stalls occur. The lower output record rate also does not demonstrate data loss. With KTable caching, multiple updates for the same key can be coalesced before forwarding. Matching final values provide the relevant correctness evidence for this aggregation.

If the hot-key latency must improve, the aggregation can be redesigned to distribute partial counts across shard keys and then combine them into the required final per-key count.

Why each option fits or fails:

A. The dominant key remains confined to one nearly saturated task, while idle threads and correct final values show limited parallelism rather than lost aggregate results.

B. More partitions create tasks, but the unchanged dominant key still routes to one partition and therefore does not parallelize most state updates.

C. A cached KTable may coalesce intermediate updates, so output record rate need not equal input rate when the final values are correct.

D. State-store operations track processing and report no write stalls, while the verified final values contradict the claim that aggregation is incomplete.


Question 16

Topic: Kafka Connect

A source connector currently writes orders.v1. The team wants registry-backed Avro values, but post-cutover sink connectors must be able to replay the complete history using one value converter.

Existing record:

topic=orders.v1 partition=0 offset=184
value.bytes.hex=7b226f726465724964223a22413137222c22746f74616c223a34322e357d
UTF-8 view={"orderId":"A17","total":42.5}

Simplified connector settings (pseudoconfiguration):

# Current
output.topic=orders.v1
value.converter=org.apache.kafka.connect.json.JsonConverter
value.converter.schemas.enable=false

# Intended
value.converter=io.confluent.connect.avro.AvroConverter
value.converter.schema.registry.url=https://registry.example.com

Assume a Kafka 4.0-compatible Connect runtime with the indicated Avro converter installed. The source can be paused and later resumed from its retained source offsets. A controlled outage is acceptable, and the post-cutover topic must contain only registry-framed Avro values.

Which transition plan best preserves record meaning and satisfies the replay requirement?

Options:

  • A. Pause the source, change its converter in place, and configure sinks to select JSON or Avro decoding according to each record’s offset.

  • B. Pause the source, decode the complete orders.v1 history as JSON into equivalent Avro records in a new topic, then resume new source output there using the Avro converter.

  • C. Pause the source, change its converter in place on orders.v1, configure sinks with the Avro converter, and replay the topic from the beginning.

  • D. Pause the source, copy the existing value bytes unchanged into a new topic, then resume new source output there using the Avro converter.

Best answer: B

Explanation: Kafka stores the bytes produced by a converter; changing converter configuration does not reinterpret or rewrite records already in a topic. The existing value begins with a JSON object and was produced by JsonConverter with schemas disabled. A registry-backed AvroConverter expects its own framing and registered schema identifier, so it cannot replay that historical value directly.

A safe homogeneous transition therefore requires semantic migration rather than a byte copy. Pause the source at a controlled boundary, read every legacy value using its JSON contract, map it to the intended Avro schema, and serialize it into a new topic with the Avro converter. After the backfill reaches the boundary, resume the source against the new topic while retaining its source offsets. Consumers can then replay the new topic from the beginning with one Avro converter.

Why each option fits or fails:

A. Offset-based dual decoding could handle a known boundary, but it does not produce the required homogeneous topic readable with one value converter.

B. The backfill preserves historical values while re-encoding them, and the controlled cutover leaves the new topic uniformly Avro for complete replay.

C. The converter change affects only future output, so the Avro converter cannot decode the existing plain JSON bytes at offset 184.

D. Copying bytes preserves the old JSON encoding, causing the new topic to contain both plain JSON and registry-framed Avro values.


Question 17

Topic: Kafka fundamentals

An order service has three independent consumers, each using a distinct consumer group:

  • Fulfillment needs accepted line items and the shipping address when OrderAccepted completes.
  • Accounting needs the payment reference and captured amount when PaymentCaptured completes.
  • Analytics needs every lifecycle transition and its occurrence time.

Records must preserve these completed business facts for one year. A new consumer must replay them without querying the order database or comparing adjacent snapshots.

The service currently writes full OrderSnapshot records with the order ID as the key to a topic configured with cleanup.policy=compact.

A replay test recorded this evidence:

offset=101 key=o-73 status=ACCEPTED updatedAt=10:00
offset=144 key=o-73 status=PAYMENT_CAPTURED paymentId=p-9 updatedAt=10:05
offset=230 key=o-73 status=DISPATCHED tracking=t-4 updatedAt=10:20
After log cleaning: only offset 230 remains
Fresh analytics group: received only offset 230
Consumer errors: none

Which design change is best supported by the evidence and requirements?

Options:

  • A. Publish one OrderCompleted record after dispatch to a one-year event topic, including the acceptance, payment, and shipment data in the final payload.

  • B. Publish a full OrderSnapshot after every database update to a one-year event topic keyed by order ID, requiring consumers to infer business transitions by comparing consecutive versions.

  • C. Publish immutable OrderAccepted, PaymentCaptured, and ShipmentDispatched records when each fact completes to a one-year event topic keyed by order ID, while keeping compacted snapshots in a separate current-state topic.

  • D. Publish separate StartFulfillment, RecordPayment, and UpdateAnalytics records at each workflow stage, using a dedicated topic for each existing consumer.

Best answer: C

Explanation: Log compaction supports a current-state representation: for records sharing a key, older values can be removed after cleaning. The trace therefore shows no consumer-group failure; the historical snapshots are absent from the cleaned log.

Business events should instead be emitted when meaningful facts complete. OrderAccepted contains the accepted lines and shipping address known at acceptance, PaymentCaptured contains the confirmed payment facts, and ShipmentDispatched contains carrier details. Keying these records by order ID supports per-order partition ordering under stable routing. An event topic retained for at least one year provides replayable history for independent consumer groups. A separate compacted snapshot topic can still expose the latest order state, but it should not serve as the required lifecycle history.

Why each option fits or fails:

A. A terminal summary is produced too late for consumers that must react when acceptance and payment facts individually complete.

B. Row-level snapshots make business boundaries implicit and require the adjacent-record comparison that the replay requirement specifically excludes.

C. Completed business-event boundaries preserve independently meaningful history, while the separate compacted projection continues to provide efficient current-state lookup.

D. Consumer-specific commands couple the record model to current applications instead of preserving shared business facts for independent and future consumers.


Question 18

Topic: Application testing

A Java application uses Apache Kafka 4.0 brokers and Java clients in a read-process-write flow.

Configuration:

  • The input consumer uses enable.auto.commit=false and isolation.level=read_committed.
  • The producer uses transactional.id=orders-worker-0 with idempotence-compatible settings.
  • For each batch, the application calls beginTransaction(), sends output records, calls sendOffsetsToTransaction() with the next input offsets, and then calls commitTransaction().
  • A separate output consumer uses isolation.level=read_committed.
  • The test inspects committed input offsets through the Admin API.

Failure cases:

  • The transaction is paused after its output send completes and its input offsets are staged.
  • An abortable transaction failure is injected, followed by a successful abortTransaction().
  • commitTransaction() throws TimeoutException, leaving its completion status uncertain.
  • A replacement producer initializes with the same transactional ID and fences the original producer.

Which proposed integration-test expectation or recovery action is INCORRECT?

Options:

  • A. The test should verify that the open transaction’s output remains unavailable to the read-committed consumer and the committed input offset remains unchanged.

  • B. The test should call abortTransaction() immediately after the commit timeout and then reprocess the batch from the prior committed offset.

  • C. The test should close the fenced producer without attempting an abort and allow only the replacement producer to continue using the transactional ID.

  • D. The test should verify that successfully aborted output remains unavailable to the read-committed consumer and the prior committed input offset is retained.

Best answer: B

Explanation: A Kafka transaction atomically exposes the produced Kafka records and the consumed offsets passed to sendOffsetsToTransaction(). A read_committed consumer cannot return records from an open or aborted transaction. Likewise, staged input offsets do not become the consumer group’s committed offsets unless the transaction commits.

A successful abort therefore leaves both the output and offset progress invisible. Fencing is different: the original producer has lost ownership of its transactional identity and must close.

A TimeoutException from commitTransaction() is especially important in recovery testing because it does not prove that the transaction aborted. The commit might have completed or might still complete. The application must not switch immediately to abortTransaction() and reprocess the batch. It should follow the supported uncertain-commit recovery path, such as retrying the commit where appropriate or closing and reconciling the outcome. These guarantees cover Kafka output records and Kafka consumer offsets, not arbitrary external side effects.

Why each option fits or fails:

A. Output records and staged offsets become visible atomically only when the open transaction commits.

B. A commit timeout leaves the outcome uncertain, so switching to abort could conflict with a commit that has already completed or is still completing.

C. Fencing is fatal for the original producer, which must close rather than attempt further transactional operations.

D. Aborting discards the transaction’s visibility and prevents its staged consumer offsets from being committed.


Question 19

Topic: Kafka Connect

A team captures customer and order changes with Kafka Connect source connectors.

Pipeline facts:

  • customers is keyed by customerId.
  • orders is keyed by orderId; each value contains customerId.
  • Only standard Kafka Connect single-record transforms are available, and neither connector has a built-in join capability.
  • A Java Kafka Streams application may be placed before the sink.

Requirement: Each new order must be enriched with the latest customer tier already processed when that order arrives. Customer updates need not re-emit earlier orders, and external lookups are prohibited.

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

Text description

Customer c9 with tier Gold passes through a source connector to the customers topic as event 1. Order o17 for customer c9 passes through another source connector to the orders topic as event 2. Both topics can supply a transformation stage whose enriched output goes to the sink connector.

Which design meets the transformation requirement?

Options:

  • A. Use a source-connector transform chain that rekeys orders by customer ID, reads the customers topic, and adds the matching tier before publication.

  • B. Use a sink-connector transform chain that consumes both topics, correlates records by customer ID, and writes each enriched order to the target.

  • C. Use log compaction on the customers topic, then let the sink correlate each order with the remaining customer record by customer ID.

  • D. Use a Kafka Streams stage that reads customers as a KTable, rekeys orders by customer ID, joins the order KStream, and writes an enriched topic.

Best answer: D

Explanation: Kafka Connect single-record transforms perform stateless reshaping of individual records, such as renaming fields, filtering, or changing keys. They do not perform joins or maintain a cross-record view of another topic.

A Kafka Streams application provides the required stateful boundary. It can materialize customers as a KTable keyed by customerId. Because orders initially use orderId as their key, the application rekeys them by customerId before joining the order KStream with the customer KTable. Since the customer event has already been processed, order o17 joins with customer c9 and receives tier Gold. A KStream-KTable join emits when an order arrives; a later customer update does not automatically re-emit previous orders, matching the stated requirement.

Why each option fits or fails:

A. A source transform can reshape or rekey its current record, but it cannot consume another topic and maintain cross-record customer state.

B. Standard sink transforms process individual records; subscribing to both topics does not provide a keyed join or customer-state lookup.

C. Compaction eventually retains newer values per key within a topic, but it neither merges topics nor supplies the sink with correlation logic.

D. The KTable holds the latest customer value, while each rekeyed order triggers a lookup and emits the required enriched record.


Question 20

Topic: Application development

A Java consumer reads a shared topic whose values use Confluent Avro framing: one magic byte, a four-byte schema ID, and the Avro payload.

Documented contract:

  • Each record must have exactly one event-type header and one event-version header.
  • The header pair must match this deployed schema allowlist:
    • Schema ID 31: OrderCreated, version 1
    • Schema ID 48: OrderCreated, version 2
    • Schema ID 62: CustomerUpdated, version 1
  • Records with missing, duplicate, unsupported, or contradictory metadata must be rejected before handler invocation.

The following language-neutral pseudocode represents the current consumer:

type = utf8(record.headers.lastHeader("event-type").value)
version = parseInt(record.headers.lastHeader("event-version").value)
reader = latestReaders[type]
event = reader.decode(record.value[5..])
handlers[type].handle(event)

A received record has headers event-type=OrderCreated and event-version=2, but its value contains schema ID 62.

Which consumer change correctly enforces the contract and prevents ambiguous decoding?

Options:

  • A. Resolve the schema ID to the allowlisted type and version, require exactly one matching header of each kind, and reject any mismatch before dispatch.

  • B. Select the latest reader by event type, use Avro compatibility to interpret any version, and reject only records whose event type is unsupported.

  • C. Resolve the schema ID and dispatch using the decoded Avro record name, while treating the type and version headers as informational metadata.

  • D. Select the reader from the exact header pair, verify that both headers occur once, and ignore the schema ID because the headers define the event contract.

Best answer: A

Explanation: Kafka stores headers and value bytes but does not enforce agreement between them. With Confluent Avro framing, the embedded schema ID identifies the writer schema actually used to encode the payload. The consumer should first validate the framing and resolve schema ID 62 through its deployed allowlist, producing CustomerUpdated version 1. It must then confirm that exactly one event-type and one event-version header are present and that both match this result. Because the supplied headers claim OrderCreated version 2, the record must be rejected before decoding or invoking a handler. Schema compatibility can support evolution within an event contract, but it does not authorize contradictory routing metadata or convert one event type into another.

Why each option fits or fails:

A. Schema ID 62 identifies CustomerUpdated version 1, so validating both headers against that mapping rejects the contradictory record before decoding or dispatch.

B. Compatibility does not make an unrelated CustomerUpdated writer schema an OrderCreated event, and the required event-version validation remains absent.

C. Schema-based dispatch avoids the wrong reader, but accepting contradictory headers violates the documented requirement that both metadata sources agree.

D. Kafka headers do not determine the Avro writer schema, so ignoring schema ID 62 could decode its payload using the incompatible OrderCreated definition.


Question 21

Topic: Application development

A Java consumer writes each order to an external ledger. The ledger operation is non-idempotent and is not part of a Kafka transaction.

Consumer settings:

  • enable.auto.commit=false
  • One consumer owns partition orders-0.
  • The application calls commitSync() only after the ledger confirms success.

Observed evidence:

Broker trace: orders-0 contains one record for order-7 at offset 41
Committed position before poll: 41
Ledger response for order-7: success
Application crash: before commitSync()
Committed position after restart: 41
Ledger response for order-7 after restart: success
Committed position after second processing: 42
Ledger rows for order-7: 2

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

Text description

The consumer starts from committed position 41, processes offset 41 in the external ledger, and crashes before committing. After restart at position 41, it processes the same record again and then commits position 42.

Which diagnostic conclusion is best supported by the evidence?

Options:

  • A. The application exhibited at-most-once processing at the ledger boundary because the ledger operation completed before the offset commit.

  • B. The producer created duplicate Kafka records, and the consumer processed each stored record once at the ledger boundary.

  • C. The manual offset commit provided exactly-once processing by atomically coordinating the ledger write with the consumer position.

  • D. The application exhibited at-least-once processing at the ledger boundary because the uncommitted Kafka record was processed again after restart.

Best answer: D

Explanation: A committed consumer offset identifies the next position to read. Before the crash, the committed position remained 41 even though processing of offset 41 had completed in the external ledger. After restart, the consumer therefore fetched offset 41 again, repeated the ledger operation, and then committed position 42.

This consume-process-commit sequence provides at-least-once processing: completed work may be repeated when a crash occurs between processing and the offset commit. The duplicate exists at the external ledger boundary; Kafka still contains only one source record. Kafka cannot roll back or deduplicate an unrelated external operation.

By contrast, committing before processing creates at-most-once behavior because a crash after the commit but before processing can cause the record’s work to be lost.

Why each option fits or fails:

A. Processing before committing creates replay risk, whereas at-most-once behavior requires committing before processing and risks loss after a crash.

B. The broker trace shows only one Kafka record at offset 41, while the unchanged committed position explains its second delivery.

C. A normal consumer offset commit is not atomic with an external ledger operation, so it cannot prevent repeating that side effect.

D. The ledger succeeded before the crash, but offset 42 was not committed, so restart replayed offset 41 and repeated the external effect.


Question 22

Topic: Application observability

A Java consumer using Apache Kafka 4.0 is configured with an OrderV2Deserializer and manual offset commits.

Baseline evidence:

  • Every restart reaches topic orders, partition 2, offset 928 and throws RecordDeserializationException.
  • The cause reports format version order-v1; the record header also contains format=order-v1.
  • The raw bytes are confirmed to be a valid Order V1 record, and a compatible V1 decoder is available.
  • Offset 928 remains the committed next position, so restarting with the same configuration repeats the failure.
  • The current handler uses the exception’s raw key/value buffers, headers, timestamp, topic, partition and offset to create a durable quarantine entry before advancing to offset 929.

The baseline completeness policy permits a record to be omitted from normal processing if that quarantine evidence is retained.

Changed condition: The policy now requires every source offset to complete normal processing; quarantine and omission are no longer sufficient. All other facts remain unchanged.

Which revised application response is required?

Options:

  • A. Install a header-aware V1/V2 deserializer, resume at offset 928, and commit offset 929 only after that record completes normal processing.

  • B. Set auto.offset.reset to earliest, restart the consumer, and rely on replay from an earlier position to bypass the failure.

  • C. Restart the consumer at offset 928 with the existing V2 deserializer until a poll eventually returns the record successfully.

  • D. Retain the quarantine handler, advance to offset 929, and treat the captured raw bytes and metadata as complete processing evidence.

Best answer: A

Explanation: A deserialization failure at a fixed partition and offset is often a poison-pill condition, not a transient consumer membership problem. Restarting resumes from the same committed position and encounters the same incompatible bytes.

Kafka 4.0 RecordDeserializationException exposes the failed record’s raw key/value buffers, headers and timestamp. Those fields support a traceable quarantine-and-skip policy, but preserving evidence is different from completing normal processing.

Under the changed policy, the application must make offset 928 readable and process it. A header-aware deserializer can route order-v1 records to the available V1 decoder while retaining V2 handling for newer records. The consumer should resume at offset 928 and advance the committed position to 929 only after processing succeeds. Changing auto.offset.reset is irrelevant because the committed offset is valid rather than missing or out of range.

Why each option fits or fails:

A. The compatible V1 path repairs the format mismatch and ensures offset 928 is processed before the committed next position advances.

B. Offset reset is not applied while a valid committed offset exists, and replay would not make offset 928 compatible with the V2 deserializer.

C. Restarting does not change the incompatible record bytes or deserializer, so the same offset continues to fail deterministically.

D. The quarantine entry preserves traceability but does not satisfy the changed requirement that offset 928 complete normal processing.


Question 23

Topic: Kafka Streams

Using Apache Kafka 4.0 Java Streams, a developer builds this stateless topology:

record Order(String status, long grossCents) {}
record Receipt(String receiptId, long netCents) {}

KStream<String, Order> input = builder.stream("orders");

input
    .filter((key, order) ->
        order != null && "PAID".equals(order.status()))
    .mapValues((readOnlyKey, order) ->
        new Receipt(
            readOnlyKey + "-" + order.grossCents(),
            order.grossCents() - 1_000L))
    .filter((key, receipt) -> receipt.netCents() >= 5_000L)
    .to("receipts");

The records enter the topology with these timestamps:

  • R1: key acct-A, value Order("PAID", 7000), timestamp 1000
  • R2: key acct-B, value Order("PAID", 5500), timestamp 2000
  • R3: key acct-C, value Order("PENDING", 9000), timestamp 3000
  • R4: key acct-D, value Order("PAID", 6000), timestamp 4000
  • R5: key acct-E, value Order("PAID", 8000), timestamp 5000

Assume serialization succeeds and no custom processor changes timestamps. Which TWO conclusions about records written to receipts are correct?

Options:

  • A. R4 produces key acct-D, value Receipt("acct-D-6000", 5000), and timestamp 4000.

  • B. R5 produces key acct-E-8000, value Receipt("acct-E-8000", 7000), and timestamp 5000.

  • C. R3 produces key acct-C, value Receipt("acct-C-9000", 8000), and timestamp 3000.

  • D. R1 produces key acct-A, value Receipt("acct-A-7000", 6000), and timestamp 1000.

  • E. R2 produces key acct-B, value Receipt("acct-B-5500", 4500), and timestamp 2000.

Correct answers: A and D

Explanation: The first filter retains only paid orders, removing R3 before transformation. mapValues then subtracts 1,000 cents and constructs a Receipt; it does not change the record key or timestamp.

The resulting net amounts are 6,000 for R1, 4,500 for R2, 5,000 for R4, and 7,000 for R5. The second filter removes R2 because 4,500 is below the threshold. Its inclusive comparison retains R4 at exactly 5,000.

Therefore, R1, R4, and R5 reach the output topic. R5 retains key acct-E, rather than adopting the receipt identifier as its key.

Why each option fits or fails:

A. The paid order maps to exactly 5,000 net cents, which satisfies the inclusive >= 5_000L predicate.

B. The mapping function may read the key when constructing the value, but mapValues leaves the record key as acct-E.

C. The PENDING status fails the first predicate, so the mapping function is never applied.

D. Both predicates pass, while mapValues calculates the new value without changing the original key or timestamp.

E. The transformed net amount is 4,500 cents, so the second filter removes the record.


Question 24

Topic: Kafka Streams

An engineer deploys an Apache Kafka Streams 4.0 application.

Fixed deployment facts:

  • bootstrap.servers lists reachable brokers.
  • security.protocol=SASL_SSL and sasl.mechanism=SCRAM-SHA-512 are set.
  • Valid SCRAM credentials and a trusted TLS truststore are supplied.
  • application.id=payment-tier-v1.
  • The default key Serde is Serdes.StringSerde.
  • The default value Serde is Serdes.LongSerde.
  • The topics and required authorizations already exist.

Baseline topology:

KStream<String, Long> payments = builder.stream(
    "payments", Consumed.with(Serdes.String(), Serdes.Long()));
KStream<String, String> tiers = payments.mapValues(
    value -> value >= 100L ? "HIGH" : "STANDARD");
tiers.to("payment-tiers",
    Produced.with(Serdes.String(), Serdes.String()));

The baseline deployment successfully processes records. In a changed deployment, the engineer changes only the final operation to:

tiers.to("payment-tiers");

What consequence occurs when the changed deployment processes its first output record?

Options:

  • A. It starts and processes successfully because the KStream<String, String> generic type automatically selects the String value Serde.

  • B. It starts and connects, then fails while deserializing the source because removing Produced also removes the source’s explicit Long Serde.

  • C. It starts and connects, then fails while serializing the String output value with the configured default Long value Serde.

  • D. It fails while constructing the topology because every sink following mapValues must declare a Produced configuration.

Best answer: C

Explanation: Kafka Streams selects runtime Serdes from explicit operator configurations or configured defaults, not from Java generic type parameters. The source still uses the explicit String key and Long value Serdes supplied through Consumed.with.

After mapValues, the value type changes from Long to String. The baseline therefore needs the String value Serde specified through Produced.with at the output boundary. Removing that override is syntactically valid and does not affect broker connectivity, security, or application.id. However, the sink falls back to the configured default Long value Serde. When it receives "HIGH" or "STANDARD", serialization fails because the Long serializer cannot serialize a String.

The unchanged application ID continues to identify the Streams application and namespace its consumer-group and internal resources.

Why each option fits or fails:

A. Java generic types do not cause Kafka Streams to discover or instantiate a runtime Serde.

B. The source remains governed by its unchanged Consumed.with configuration; downstream sink settings do not alter source deserialization.

C. Removing the sink override leaves the String-valued stream subject to the incompatible default Long serializer at the output boundary.

D. The overload without Produced is valid, so the incompatibility appears during record serialization rather than topology construction.


Continue in the web app

Use IT Mastery for interactive Confluent CCDAK practice with mixed sets, timed mocks, topic drills, explanations, and progress tracking.

Try Confluent CCDAK on Web

Recall the key distinctions · Official resources · Report a question issue