End to end: load JSON, Columnar and Binary Formats's 100,000 orders into booknest.orders, start App, wait until its totals stop changing, and compare them with sums computed straight from the files by streams/expected.py:
# Load the orders, run the application until its totals settle, check them against the files
K="docker exec -i l3-kafka /opt/kafka/bin"; B="--bootstrap-server l3-kafka:9092"
$K/kafka-topics.sh $B --create --topic booknest.orders --partitions 3 >/dev/null
jq -rc '"\(.order_id)|\(.)"' data/orders.jsonl | $K/kafka-console-producer.sh $B \
--topic booknest.orders --reader-property parse.key=true --reader-property key.separator='|'
cd streams && start=$SECONDS
java -cp 'lib/*:classes' booknest.App 2> app.log | python3 expected.py
echo "settled after $(( SECONDS - start )) s"
echo "late lines skipped by the daily windows: $(grep -c 'expired window' app.log)"Output
Cooking 716,952.00 716,952.00 Fiction 693,362.45 693,362.45 Home and Garden 334,985.10 334,985.10 Science Fiction 533,466.00 533,466.00 Technology 934,530.50 934,530.50 Travel 90,131.25 90,131.25 all equal: True settled after 177 s late lines skipped by the daily windows: 85149
The totals match the files to the cent, but the daily windows skipped 85,149 order lines as later than their one-hour grace: replaying history, stream time races ahead when partitions are read at different speeds, so widen the grace for backfills. App runs at least once; under exactly_once_v2 this busy host needed producer.transaction.timeout.ms above Streams' 10-second default, or transactions aborted (InvalidProducerEpochException).