From 58a5866afa8eaff176fc22639f8124fb8a966293 Mon Sep 17 00:00:00 2001 From: Mark Reitblatt Date: Tue, 18 Aug 2026 11:56:21 -0700 Subject: [PATCH] Fix excessive CPU load in generate_test_data.sh (#455) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * generate_test_data.sh was spawning a brand-new kafka-console-producer JVM process per message (once/sec), paying full JVM startup + classpath scan + producer init/metadata fetch every iteration, forever — this pegged the osprey-kafka-test-data-producer container at 180-200% CPU cycling every 1-2s under docker top. * Now starts a single long-lived kafka-console-producer process and streams generated messages to it over a FIFO, instead of forking a new producer per line. * Added 'unless-stopped' restart policy to test data generator container Co-authored-by: Claude Sonnet 5 --- docker-compose.yaml | 1 + example_data/generate_test_data.sh | 45 ++++++++++++++++++++++++------ 2 files changed, 38 insertions(+), 8 deletions(-) diff --git a/docker-compose.yaml b/docker-compose.yaml index 53a167f..6ec8ac5 100644 --- a/docker-compose.yaml +++ b/docker-compose.yaml @@ -222,6 +222,7 @@ services: entrypoint: - /bin/bash command: ["/osprey/example_data/generate_test_data.sh"] + restart: unless-stopped postgres: hostname: postgres diff --git a/example_data/generate_test_data.sh b/example_data/generate_test_data.sh index 13b6283..10a0723 100755 --- a/example_data/generate_test_data.sh +++ b/example_data/generate_test_data.sh @@ -45,21 +45,40 @@ generate_action() { eval "$cmd" "$SCRIPT_DIR/template.json" } -# Function to build kafka-console-producer command +# Function to build kafka-console-producer command as an array (avoids +# eval'ing environment-controlled values as shell code) build_kafka_command() { - local cmd="kafka-console-producer --broker-list $KAFKA_BROKER --topic $KAFKA_TOPIC" + kafka_cmd=(kafka-console-producer --broker-list "$KAFKA_BROKER" --topic "$KAFKA_TOPIC") if [ -n "$KAFKA_CONFIG_FILE" ]; then - cmd="$cmd --producer.config $KAFKA_CONFIG_FILE" + kafka_cmd+=(--producer.config "$KAFKA_CONFIG_FILE") fi - - echo "$cmd" } +# Directory + FIFO used to stream messages into a single long-lived producer +# process. Using our own mktemp -d directory avoids writing into shared /tmp. +FIFO_DIR="$(mktemp -d)" +FIFO="$FIFO_DIR/kafka_producer_fifo" +producer_pid="" + # Function to handle cleanup on script termination cleanup() { echo echo "Stopping data generation..." + # Close our write end of the FIFO so the producer sees EOF and exits + exec 3>&- 2>/dev/null + if [ -n "$producer_pid" ]; then + # Give the producer a chance to exit on its own after seeing EOF + for _ in $(seq 1 20); do + kill -0 "$producer_pid" 2>/dev/null || break + sleep 0.1 + done + if kill -0 "$producer_pid" 2>/dev/null; then + kill "$producer_pid" 2>/dev/null + fi + wait "$producer_pid" 2>/dev/null + fi + rm -rf "$FIFO_DIR" exit 0 } @@ -75,14 +94,24 @@ fi echo "Press Ctrl+C to stop..." echo -# Build the kafka command -kafka_cmd=$(build_kafka_command) +# Build the kafka command (populates the kafka_cmd array) +build_kafka_command + +# Start a single long-lived producer reading from the FIFO, instead of +# spawning a brand-new JVM producer process per message. +mkfifo "$FIFO" || { echo "Failed to create FIFO at $FIFO" >&2; exit 1; } +"${kafka_cmd[@]}" < "$FIFO" & +producer_pid=$! + +# Open the FIFO for writing on fd 3 and keep it open for the life of the +# script, so the producer never sees EOF between messages. +exec 3> "$FIFO" # Infinite loop to generate and send actions while true; do action=$(generate_action) echo -e "Sending $action" - echo -e "$action" | $kafka_cmd + echo -e "$action" >&3 # Increment action_id in the main shell ((action_id++)) -- 2.51.2