Prev Next

Integration / Apache Kafka Interview questions

Explain the execution flow of a Kafka Streams topology processing a record?

A Kafka Streams application defines a topology — a graph of processing nodes (source, transformations, sink) — and every record pulled from an input topic flows through that graph node by node before any result is produced.

flowchart TD A[Record polled from source topic partition] --> B[Deserialized into a typed key/value pair] B --> C[Passed to the first processor node in the topology] C --> D{Node type} D -- Stateless, e.g. filter/map --> E[Transform or drop the record, pass downstream] D -- Stateful, e.g. aggregate/join --> F[Read/update local state store, emit result] E --> G{More nodes downstream?} F --> G G -- Yes --> C G -- No --> H[Final result serialized and written to output topic, if a sink is defined] H --> I[If stateful, changelog topic updated to reflect the state change]

This per-record flow happens continuously as the application polls its input topics, and because the topology is defined once at startup (via the DSL's fluent builder or the lower-level Processor API), Kafka Streams can statically determine which internal topics (repartition, changelog) it needs to create before any data actually starts flowing, rather than discovering that requirement dynamically at runtime.

What must happen before a stateful operation like aggregate can emit a result for a record?
When does Kafka Streams determine which internal topics it needs to create?

Invest now in Acorns!!! 🚀 Join Acorns and get your $5 bonus!
Acorns Logo

Invest now in Acorns!!! 🚀
Join Acorns and get your $5 bonus!

Earn passively and while sleeping

Acorns is a micro-investing app that automatically invests your "spare change" from daily purchases into diversified, expert-built portfolios of ETFs. It is designed for beginners, allowing you to start investing with as little as $5. The service automates saving and investing. Disclosure: I may receive a referral bonus.

Robinhood Logo

Invest now!!! Get Free equity stock (US, UK only)!

Use Robinhood app to invest in stocks. It is safe and secure. Use the Referral link to claim your free stock when you sign up!.

The Robinhood app makes it easy to trade stocks, crypto and more.


Webull Logo

Webull! Receive free stock by signing up using the link: Webull signup.

More Related questions...

What is a Kafka topic? Mention some of the Apache Kafka terminologies. What is the Global Unique identifier of a Kafka Message? What is a partition in Kafka? What is a Kafka producer? What is Apache Kafka? Explain the role of the ZooKeeper in Kafka. What is a Kafka consumer? What is a consumer group in Kafka? What is Apache ZooKeeper? What is a Kafka broker? Difference between Apache Kafka and Confluent Kafka. Kafka's zero-copy principle. Define replication in Kafka? What is KRaft mode in Kafka? Explain sequential I/O principle in Kafka. What are the types of message delivery semantics in Kafka? What is a Kafka offset used for? List the core APIs provided by Kafka? What is the purpose of a Kafka topic's retention policy? Describe the role of a partition leader in Kafka? What is a Kafka Connect connector? What are Kafka Streams used for? How do you create a topic using the Kafka CLI? What is a serializer in Kafka producer configuration? How do you list existing topics in a Kafka cluster? Why is Kafka better suited than a traditional message queue for high-throughput event streaming? How does Kafka differ from RabbitMQ? What is the difference between at-least-once, at-most-once, and exactly-once delivery semantics? How does Kafka's KRaft controller quorum manage cluster metadata? Why should you configure min.insync.replicas alongside acks=all? How does Kafka handle partition leader election? When should you increase the number of partitions for a topic? What happens when a consumer in a group fails to send a heartbeat in time? Explain the execution flow of a Kafka producer sending a message to a broker? How can you optimize Kafka producer throughput? How do you troubleshoot consumer lag in a Kafka application? Why is the in-sync replica (ISR) set important for durability? Explain the lifecycle of a Kafka consumer group rebalance? How does Kafka achieve exactly-once semantics with idempotent producers and transactions? What is the difference between log compaction and log deletion cleanup policies? How do you implement a custom partitioner in Kafka? How does Kafka Streams manage local state using state stores? Which is better and why: the classic consumer rebalance protocol or the new KIP-848 protocol? How do you integrate a schema registry with Kafka producers and consumers? Explain the internal working of Kafka's replication protocol between leader and follower brokers? How do you configure tiered storage for a Kafka topic? What is the difference between Kafka Connect source and sink connectors? How does Kafka support message compression? When would you choose the cooperative sticky partition assignment strategy? How do you secure a Kafka cluster with SASL and ACLs? Why should unclean leader election be disabled in most production clusters? How do you configure MirrorMaker for cross-cluster replication? Explain the execution flow of a Kafka Streams topology processing a record? How does Kafka report and expose broker and consumer metrics for monitoring? Why doesn't increasing partition count always improve throughput? How do you migrate a Kafka cluster from ZooKeeper mode to KRaft mode? What is the difference between Kafka's Queues feature (KIP-932) and traditional partitioned consumption?
Show more question and Answers...

Akka interview questions

Comments & Discussions