Learn Labs
7. Reliable Data Delivery

7.6 Validating system reliability

Note the shape of that expectation: a quantified duplicate budget.

"Once we have gone through the process of figuring out our reliability requirements, configuring the brokers, configuring the clients, and using the APIs in the best way for our use case, we can just relax and run everything in production, confident that no event will ever be missed, right?"

Three layers of validation:

  1. Validate the Configuration
  2. Validate the Application
  3. Monitor the application in production

6.1 Validating configuration

Two reasons to test config in isolation from application logic:

  • “It helps to test if the configuration we’ve chosen can meet our requirements.”
  • “It is a good exercise to reason through the expected behavior of the system.”

The tools: org.apache.kafka.tools includes VerifiableProducer and VerifiableConsumer — "These can run as command-line tools or be embedded in an automated testing framework."

produceconsumeproduced countconsumed countVerifiableProducernumbers from 1 to a value we choosetopic partitionVerifiableConsumerprints the events IN ORDERdo the counts MATCH?produced vs consumed
  • · VerifiableProducer — configured like your real producer: acks, retries, delivery.timeout.ms, and the rate at which messages are produced. It “prints SUCCESS OR ERROR for each message sent, based on the acks received.”
  • · VerifiableConsumer — “prints out the events it consumed IN ORDER” and “also prints information regarding COMMITS and REBALANCES.”
  • · The test: run a failure scenario, then check that “the number of messages PRODUCED by the producer and the number of messages CONSUMED by the consumer MATCH.”
Figure 7.6.26.1 Validating configuration

The four scenarios the book says to test:

TestThe question to answer
Leader election"what happens if we kill the leader? How long does it take the producer and consumer to start working as usual again?"
Controller election"how long does it take the system to resume after a restart of the controller?"
Rolling restart"can we restart the brokers one by one without losing any messages?"
Unclean leader election"what happens when we kill all the replicas for a partition one by one (to make sure each goes out of sync) and then start a broker that was out of sync? What needs to happen in order to resume operations? IS THIS ACCEPTABLE?"

The leader-election test is not optional trivia — its answer is the number you plug into delivery.timeout.ms (Ch. 3: "it typically takes leader election 30 seconds... let's keep retrying for 120 seconds"). You are supposed to measure it, not guess it.

"The Apache Kafka source repository includes an extensive test suite. Many of the tests are based on the same principle and use the verifiable producer and consumer to make sure rolling upgrades work."

6.2 Validating applications

"This will check things like custom error-handling code, offset commits, and rebalance listeners and similar places where the application logic interacts with Kafka's client libraries."

The eight failure conditions to test under:

  • Clients lose connectivity to one of the brokers
  • High latency between client and broker
  • Disk full
  • Hanging disk (also called a “brown out”) — the nastiest one: not failed, just glacial
  • Leader election
  • Rolling restart of brokers
  • Rolling restart of consumers
  • Rolling restart of producers

Fault injection: "Apache Kafka itself includes the Trogdor test framework for fault injection."

The methodology — write the expectation first:

"For each scenario, we will have EXPECTED BEHAVIOR, which is what we planned on seeing when we developed the application. Then we run the test to see what ACTUALLY happens."

Example: "when planning for a rolling restart of consumers, we planned for a short pause as consumers rebalance and then continue consumption with NO MORE THAN 1,000 DUPLICATE VALUES. Our test will show whether the way the application commits offsets and handles rebalances actually works this way."

Note the shape of that expectation: a quantified duplicate budget. "No duplicates" is not a testable expectation for an at-least-once system; "≤1,000 duplicates" is.

6.3 Monitoring reliability in production

"Testing the application is important, but it does not replace the need to continuously monitor production systems to make sure data is flowing as expected."

Producer-side metrics

"For the producers, the two metrics most important for reliability are ERROR-RATE and RETRY-RATE PER RECORD (aggregated). Keep an eye on those, since error or retry rates going up can indicate an issue."

The log lines to watch, and how to read them:

WARN level, while sending events:

Got error produce response with correlation id 5689 on topic-partition
[topic-1,3], retrying (two attempts left). Error: ...
                                 ▲▲▲▲▲▲▲▲▲▲▲▲▲▲▲▲▲

“When we see events with 0 attempts left, the producer is running out of retries.”

ERROR level is “likely to indicate that sending the message failed completely due to

  • a nonretriable error,
  • a retriable error that ran out of retries, or
  • a timeout.

When applicable, the exact error from the broker will be logged.”

(Note the correlation id in the log line — that's the request header field from Ch. 6 §5, and it's how you match this WARN to a broker-side log entry.)

"Of course, it is always better to solve the problem that caused the errors in the first place" — rather than just extending delivery.timeout.ms.

Consumer-side: lag

"the most important metric is CONSUMER LAG... Ideally, the lag would always be zero. In practice, because calling poll() returns multiple messages and then the consumer spends time processing them before fetching more, the lag will always fluctuate a bit. What is important is to make sure consumers do EVENTUALLY CATCH UP rather than fall further and further behind."

"Because of the expected fluctuation in consumer lag, setting traditional alerts on the metric can be challenging. Burrow is a consumer lag checker by LinkedIn and can make this easier."

lag oscillates by design — poll returns a batch, then you processwhat you want to detect is the TREND: is it recovering?consumer lagtime

A static threshold either flaps or is set uselessly high — hence Burrow’s sliding-window, status-based evaluation.

Figure 7.6.5Consumer-side: lag
End-to-end flow monitoring

"Monitoring flow of data also means making sure all produced data is consumed in a timely manner ('timely manner' is usually based on business requirements). In order to make sure data is consumed in a timely manner, we need to know when the data was produced."

Kafka helps: "starting with version 0.10.0, all messages include a timestamp that indicates when the event was produced (although note that this can be overridden either by the application that is sending the events or by the brokers themselves if they are configured to do so)."

(That's the create-time vs append-time distinction from Ch. 6's batch header. If brokers are set to append-time, your "produce latency" measurement is measuring the wrong thing.)

What you must build:

events/secevents/secPRODUCERS recordevents produced (usually events/second)CONSUMERS recordevents consumed per unit of timeRECONCILIATION SYSTEM…and the LAG FROM PRODUCE-TIME TO CONSUME-TIME, using the event timestamp
  • · compare events/sec from producer vs consumer → “to make sure no messages were lost on the way”
  • · check the interval between produce time and consume time is reasonable
Figure 7.6.6What you must build

The honest caveat: "This type of end-to-end monitoring system can be challenging and time-consuming to implement. To the best of our knowledge, THERE IS NO OPEN SOURCE IMPLEMENTATION of this type of system, but Confluent provides a commercial implementation as part of the Confluent Control Center."

Broker-side error metrics
 kafka.server:type=BrokerTopicMetrics,name=FailedProduceRequestsPerSec
 kafka.server:type=BrokerTopicMetrics,name=FailedFetchRequestsPerSec

"At times, some level of error responses is expected — for example, if we shut down a broker for maintenance and new leaders are elected on another broker, it is expected that producers will receive a NOT_LEADER_FOR_PARTITION error, which will cause them to request updated metadata before continuing as usual. Unexplained increases in failed requests should ALWAYS be investigated. To assist in such investigations, the failed requests metrics are TAGGED WITH THE SPECIFIC ERROR RESPONSE that the broker sent."

The tagging is the important operational detail: NOT_LEADER_FOR_PARTITION during a rolling restart is expected; the same counter rising with NOT_ENOUGH_REPLICAS means your min.insync.replicas floor is being hit and producers are being rejected.


On this page