Test Yourself!

These ten questions cover the chapter, and each catches people who already run Kafka 129 : an order split across partitions, a compaction that keeps everything, a minimum ISR that is not enforced, a size limit a codec gets around, a consumer that reads nothing, records hidden by another producer's transaction, a job that vanishes, a connector that drops a record without a word, a rename the registry half accepts, and a late order a window ignores. All ran on the single broker l3-kafka (Kafka 4.3.1, Single-Broker KRaft Cluster) with BookNest's order events, confluent-kafka 2.15.1 and Schema Registry 8.3.2. Write down what each prints and why, then check Appendix I for the real output, the reason, the fix and the section.

check.sh in demos/ch06/quiz/ starts a fresh broker and registry, loads the events, creates the quiz topics and runs the three files; lib.sh and quizlib.py hold the helpers their comments name.

questions.sh: questions 1-3, 8 and 9 with the command-line toolsShell
. quiz/lib.sh       # topic NAME OPTS...; put TOPIC TYPE... (order 7, acks=all); all TOPIC
# 1. Order 7's events, two before and two after a resize. Where do they land? (6.1.3, 6.4.1)
topic booknest.q1 --create --partitions 3; put booknest.q1 order_placed order_paid
topic booknest.q1 --alter --partitions 6; put booknest.q1 order_shipped order_delivered
all booknest.q1 | jq -Rr 'split("\t") | "1: \(.[0]) \(.[1] | fromjson | .type)"' | sort -sk2,2
# 2. Order 7's four events in a compacted topic. A minute later, how many are left? (6.2.4)
topic booknest.q2 --create --config cleanup.policy=compact \
  --config min.cleanable.dirty.ratio=0.01
put booknest.q2 order_placed order_paid order_shipped order_delivered; sleep 60
echo "2: $(all booknest.q2 | wc -l) records"
# 3. One replica, min.insync.replicas=2, acks=all: is the write accepted? (6.2.6, 6.4.2)
topic booknest.q3 --create --replication-factor 1 --config min.insync.replicas=2
put booknest.q3 order_placed; echo "3: $(all booknest.q3 | wc -l) record(s)"
# 8. A FileStreamSink with errors.tolerance=all (q8-sink.properties) reads four records,
#    one not JSON, with a schemaless JSON converter. What reaches the file? (6.8.8)
printf '%s\n' '{"order_id":7}' 'not json' '{"order_id":8}' '{"order_id":9}' |
  $K/kafka-console-producer.sh $B --topic booknest.q8
docker exec l3-kafka timeout 30 /opt/kafka/bin/connect-standalone.sh /tmp/q8.properties \
  /tmp/q8-sink.properties > q8.log 2>&1
echo "8: $(docker exec l3-kafka cat /tmp/q8.txt | xargs); ERROR lines: $(grep -c ERROR q8.log)"
# 9. Rename the field type of booknest.q9-value (version 1: order_event.avsc) to event_type
#    with an alias: compatible under BACKWARD? Under FULL? (Section 6.9.6)
R=localhost:33081; S=booknest.q9-value; H='Content-Type: application/json'
for mode in BACKWARD FULL; do
  curl -s -X PUT -H "$H" -d "{\"compatibility\": \"$mode\"}" $R/config/$S >/dev/null
  jq -c '.fields[2] += {name: "event_type", aliases: ["type"]}' schemas/order_event.avsc |
    jq -Rn '{schema: input}' |
    curl -s -X POST -H "$H" -d @- $R/compatibility/subjects/$S/versions/latest |
    jq -r --arg m $mode '"9: \($m) \(.is_compatible)"'
done
quiz.py: questions 4-7 with confluent-kafkaPython
import json
from confluent_kafka import AcknowledgeType as Ack, Producer, ShareConsumer
from quizlib import B, console, count, send   # B = bootstrap; send gives "ok" or the error
# 4. A nightly export, 6,000 events in one message, goes to booknest.q4, whose
#    max.message.bytes is 262144: uncompressed, then with zstd. (Sections 6.4.2, 6.5.5)
export = json.dumps([json.loads(line) for line in open("data/order_events.jsonl")][:6000])
print("4:", len(export), [send("booknest.q4", export, **{"compression.type": codec})
                          for codec in ("none", "zstd")])
# 5. How many records does a new group read from the full topic in 10 seconds? (6.7.1)
print("5:", count("booknest.order-events"))
# 6. A transaction stays open while three plain records follow it into booknest.q6 (one
#    partition). What do a Python and the console consumer read? (Section 6.6.3)
tx = Producer({**B, "transactional.id": "booknest-q6"})
tx.init_transactions(); tx.begin_transaction()
tx.produce("booknest.q6", "refund 1"); tx.flush()          # sent, not committed
for n in (2, 3, 4): send("booknest.q6", f"refund {n}")
print("6:", count("booknest.q6", **{"auto.offset.reset": "earliest"}), console("booknest.q6"))
tx.commit_transaction()
# 7. booknest.q7 holds jobs 1-3; a packer releases job 3 whenever it gets it. (6.7.10)
c = ShareConsumer({**B, "group.id": "booknest-q7", "share.acknowledgement.mode": "explicit"})
c.subscribe(["booknest.q7"]); got = {}
for _ in range(15):
    for m in c.poll(1.0):
        got[m.value().decode()] = got.get(m.value().decode(), 0) + 1
        c.acknowledge(m, Ack.RELEASE if m.value() == b"job 3" else Ack.ACCEPT)
    c.commit_sync()
c.close(); print("7:", got)
Quiz.java: question 10 with Kafka StreamsJava
import java.time.*;
import org.apache.kafka.common.serialization.Serdes;
import org.apache.kafka.streams.*;
import org.apache.kafka.streams.kstream.*;
import org.apache.kafka.streams.state.QueryableStoreTypes;
public class Quiz {
  // 10. booknest.q10 holds three Fiction orders stamped 23:00 on 1 June, 02:00 on 2 June and
  //     23:30 on 1 June. Daily windows with an hour of grace: 1 June's count? (6.10.3, 6.10.7)
  public static void main(String[] args) throws Exception {
    var builder = new StreamsBuilder();
    builder.stream("booknest.q10", Consumed.with(Serdes.String(), Serdes.String()))
        .groupByKey().windowedBy(TimeWindows.ofSizeAndGrace(Duration.ofDays(1),
                                                            Duration.ofHours(1)))
        .count(Materialized.as("daily"));
    var props = new java.util.Properties();
    props.put("bootstrap.servers", "localhost:33092");
    props.put("application.id", "booknest-q10-" + System.currentTimeMillis());
    try (var streams = new KafkaStreams(builder.build(), props)) {
      streams.start(); Thread.sleep(20_000);
      try (var days = streams.store(StoreQueryParameters.fromNameAndType("daily",
          QueryableStoreTypes.<String, Long>windowStore())).all()) {
        days.forEachRemaining(d -> System.out.println("10: " + d.key.window().startTime()
                                                      + " " + d.value));
      }
    }
  }
}