Configure Kafka Streams state stores with RocksDB optimization for high-performance streaming applications. Learn custom state store configurations, RocksDB tuning parameters, and monitoring techniques for production-grade stream processing.
Prerequisites
- Java 11 or higher installed
- Apache Kafka cluster running
- Maven build tool
- Minimum 4GB RAM available
What this solves
Kafka Streams applications rely on state stores for maintaining intermediate data during stream processing operations like aggregations, joins, and windowing. The default RocksDB state store configuration often becomes a performance bottleneck under high-throughput workloads. This tutorial shows you how to configure custom state stores, optimize RocksDB parameters, and monitor performance metrics to achieve production-grade streaming performance with lower latency and higher throughput.
Prerequisites and system setup
Update system packages
Start by updating your package manager to ensure you get the latest versions of all dependencies.
sudo apt update && sudo apt upgrade -ysudo dnf update -yInstall Java runtime environment
Kafka Streams requires Java 11 or later. Install OpenJDK which provides excellent performance for streaming applications.
sudo apt install -y openjdk-17-jdk openjdk-17-jresudo dnf install -y java-17-openjdk java-17-openjdk-develVerify the Java installation:
java -versionDownload and install Apache Kafka
Download the latest Kafka distribution with Scala 2.13 binaries for optimal performance.
cd /opt
sudo wget https://downloads.apache.org/kafka/2.8.2/kafka_2.13-2.8.2.tgz
sudo tar -xzf kafka_2.13-2.8.2.tgz
sudo mv kafka_2.13-2.8.2 kafka
sudo chown -R $USER:$USER /opt/kafkaAdd Kafka binaries to your PATH:
export KAFKA_HOME=/opt/kafka
export PATH=$PATH:$KAFKA_HOME/binsource ~/.bashrcStart Kafka cluster
Start ZooKeeper and Kafka broker services for local development and testing.
cd /opt/kafka
# Start ZooKeeper
bin/zookeeper-server-start.sh -daemon config/zookeeper.properties
# Start Kafka broker
bin/kafka-server-start.sh -daemon config/server.propertiesVerify the cluster is running:
bin/kafka-broker-api-versions.sh --bootstrap-server localhost:9092Configure custom state store topology
Create Maven project structure
Set up a Maven project for your Kafka Streams application with RocksDB dependencies.
mkdir -p kafka-streams-optimization/src/main/java/com/example
cd kafka-streams-optimizationCreate the Maven POM file with required dependencies:
Implement custom RocksDB state store configuration
Create a custom state store configuration class that optimizes RocksDB parameters for high-performance streaming.
package com.example;
import org.apache.kafka.streams.state.RocksDBConfigSetter;
import org.rocksdb.*;
import java.util.Map;
public class OptimizedRocksDBConfig implements RocksDBConfigSetter {
@Override
public void setConfig(String storeName, Options options, MapCreate Kafka Streams application with optimized state stores
Implement a streaming application that uses the custom RocksDB configuration for state stores.
package com.example;
import org.apache.kafka.common.serialization.Serdes;
import org.apache.kafka.streams.KafkaStreams;
import org.apache.kafka.streams.StreamsBuilder;
import org.apache.kafka.streams.StreamsConfig;
import org.apache.kafka.streams.kstream.*;
import org.apache.kafka.streams.state.Stores;
import java.time.Duration;
import java.util.Properties;
public class StreamsApplication {
public static void main(String[] args) {
Properties props = new Properties();
props.put(StreamsConfig.APPLICATION_ID_CONFIG, "optimized-streams-app");
props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass());
props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass());
// Configure RocksDB optimization
props.put(StreamsConfig.ROCKSDB_CONFIG_SETTER_CLASS_CONFIG, OptimizedRocksDBConfig.class);
// Optimize processing threads
props.put(StreamsConfig.NUM_STREAM_THREADS_CONFIG, Runtime.getRuntime().availableProcessors());
// Configure commit interval for performance
props.put(StreamsConfig.COMMIT_INTERVAL_MS_CONFIG, 10000);
// Enable exactly-once semantics
props.put(StreamsConfig.PROCESSING_GUARANTEE_CONFIG, StreamsConfig.EXACTLY_ONCE_V2);
StreamsBuilder builder = new StreamsBuilder();
// Example: Word count with optimized state store
KStreamCreate topics and compile application
Create the required Kafka topics and compile your streaming application.
# Create input and output topics
kafka-topics.sh --create --topic input-topic --bootstrap-server localhost:9092 --partitions 4 --replication-factor 1
kafka-topics.sh --create --topic output-topic --bootstrap-server localhost:9092 --partitions 4 --replication-factor 1
# Compile the application
mvn clean compile exec:java -Dexec.mainClass="com.example.StreamsApplication"Advanced RocksDB tuning parameters
Configure memory-based optimizations
Create an advanced RocksDB configuration that optimizes memory usage patterns for different workload types.
package com.example;
import org.apache.kafka.streams.state.RocksDBConfigSetter;
import org.rocksdb.*;
import java.util.Map;
public class AdvancedRocksDBConfig implements RocksDBConfigSetter {
private static final long BLOCK_CACHE_SIZE = 128 * 1024 * 1024L; // 128MB
private static final long WRITE_BUFFER_SIZE = 64 * 1024 * 1024L; // 64MB
private static final int MAX_WRITE_BUFFER_NUMBER = 6;
@Override
public void setConfig(String storeName, Options options, MapImplement workload-specific state store configurations
Create different state store configurations optimized for specific streaming patterns like aggregations versus joins.
package com.example;
import org.apache.kafka.streams.state.RocksDBConfigSetter;
import org.rocksdb.*;
import java.util.Map;
public class WorkloadSpecificConfig implements RocksDBConfigSetter {
@Override
public void setConfig(String storeName, Options options, MapMonitor RocksDB performance metrics
Create RocksDB metrics collection
Implement a metrics collector that exposes RocksDB statistics for monitoring and alerting.
package com.example;
import org.apache.kafka.streams.KafkaStreams;
import org.apache.kafka.streams.state.QueryableStoreTypes;
import org.apache.kafka.streams.state.ReadOnlyKeyValueStore;
import org.rocksdb.RocksDB;
import org.rocksdb.Statistics;
import org.rocksdb.TickerType;
import java.util.HashMap;
import java.util.Map;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
public class RocksDBMetricsCollector {
private final KafkaStreams streams;
private final ScheduledExecutorService scheduler;
private final MapAutomated install script
Run this to automate the entire setup
#!/usr/bin/env bash
set -euo pipefail
# Colors for output
RED='\033[0;31m'
GREEN='\033[0;32m'
YELLOW='\033[1;33m'
NC='\033[0m'
# Configuration
KAFKA_VERSION="2.8.2"
SCALA_VERSION="2.13"
INSTALL_DIR="/opt"
KAFKA_USER="kafka"
PROJECT_NAME="kafka-streams-optimization"
print_status() {
echo -e "${GREEN}[INFO]${NC} $1"
}
print_warning() {
echo -e "${YELLOW}[WARN]${NC} $1"
}
print_error() {
echo -e "${RED}[ERROR]${NC} $1"
}
usage() {
echo "Usage: $0 [OPTIONS]"
echo "Options:"
echo " -u, --user USER Kafka user (default: kafka)"
echo " -d, --dir DIR Install directory (default: /opt)"
echo " -h, --help Show this help message"
exit 1
}
cleanup_on_error() {
print_error "Installation failed. Cleaning up..."
if systemctl is-active --quiet kafka 2>/dev/null; then
systemctl stop kafka
fi
if systemctl is-active --quiet zookeeper 2>/dev/null; then
systemctl stop zookeeper
fi
if [ -d "${INSTALL_DIR}/kafka" ]; then
rm -rf "${INSTALL_DIR}/kafka"
fi
if id "$KAFKA_USER" &>/dev/null; then
userdel -r "$KAFKA_USER" 2>/dev/null || true
fi
}
trap cleanup_on_error ERR
# Parse command line arguments
while [[ $# -gt 0 ]]; do
case $1 in
-u|--user)
KAFKA_USER="$2"
shift 2
;;
-d|--dir)
INSTALL_DIR="$2"
shift 2
;;
-h|--help)
usage
;;
*)
echo "Unknown option: $1"
usage
;;
esac
done
# Check if running as root
if [ "$EUID" -ne 0 ]; then
print_error "This script must be run as root or with sudo"
exit 1
fi
# Auto-detect distribution
if [ -f /etc/os-release ]; then
. /etc/os-release
case "$ID" in
ubuntu|debian)
PKG_MGR="apt"
PKG_INSTALL="apt install -y"
PKG_UPDATE="apt update && apt upgrade -y"
JAVA_PKG="openjdk-17-jdk"
;;
almalinux|rocky|centos|rhel|ol|fedora)
PKG_MGR="dnf"
PKG_INSTALL="dnf install -y"
PKG_UPDATE="dnf update -y"
JAVA_PKG="java-17-openjdk java-17-openjdk-devel"
;;
amzn)
PKG_MGR="yum"
PKG_INSTALL="yum install -y"
PKG_UPDATE="yum update -y"
JAVA_PKG="java-17-openjdk java-17-openjdk-devel"
;;
*)
print_error "Unsupported distribution: $ID"
exit 1
;;
esac
else
print_error "Cannot detect distribution. /etc/os-release not found."
exit 1
fi
print_status "Detected distribution: $ID"
echo "[1/9] Updating system packages..."
$PKG_UPDATE
echo "[2/9] Installing Java runtime environment..."
$PKG_INSTALL $JAVA_PKG wget tar
# Verify Java installation
java -version
echo "[3/9] Creating Kafka user..."
if ! id "$KAFKA_USER" &>/dev/null; then
useradd -r -m -s /bin/bash "$KAFKA_USER"
fi
echo "[4/9] Downloading and installing Apache Kafka..."
cd "$INSTALL_DIR"
if [ ! -f "kafka_${SCALA_VERSION}-${KAFKA_VERSION}.tgz" ]; then
wget "https://downloads.apache.org/kafka/${KAFKA_VERSION}/kafka_${SCALA_VERSION}-${KAFKA_VERSION}.tgz"
fi
if [ -d "kafka" ]; then
rm -rf kafka
fi
tar -xzf "kafka_${SCALA_VERSION}-${KAFKA_VERSION}.tgz"
mv "kafka_${SCALA_VERSION}-${KAFKA_VERSION}" kafka
chown -R "$KAFKA_USER:$KAFKA_USER" "${INSTALL_DIR}/kafka"
echo "[5/9] Configuring environment variables..."
cat > /etc/profile.d/kafka.sh << EOF
export KAFKA_HOME=${INSTALL_DIR}/kafka
export PATH=\$PATH:\$KAFKA_HOME/bin
EOF
chmod 644 /etc/profile.d/kafka.sh
source /etc/profile.d/kafka.sh
echo "[6/9] Creating systemd service files..."
# ZooKeeper service
cat > /etc/systemd/system/zookeeper.service << EOF
[Unit]
Description=Apache ZooKeeper
Documentation=http://zookeeper.apache.org
Requires=network.target remote-fs.target
After=network.target remote-fs.target
[Service]
Type=forking
User=$KAFKA_USER
Group=$KAFKA_USER
Environment=JAVA_HOME=/usr/lib/jvm/java-17-openjdk
ExecStart=${INSTALL_DIR}/kafka/bin/zookeeper-server-start.sh -daemon ${INSTALL_DIR}/kafka/config/zookeeper.properties
ExecStop=${INSTALL_DIR}/kafka/bin/zookeeper-server-stop.sh
TimeoutSec=30
Restart=on-failure
[Install]
WantedBy=multi-user.target
EOF
# Kafka service
cat > /etc/systemd/system/kafka.service << EOF
[Unit]
Description=Apache Kafka
Documentation=http://kafka.apache.org
Requires=zookeeper.service
After=zookeeper.service
[Service]
Type=forking
User=$KAFKA_USER
Group=$KAFKA_USER
Environment=JAVA_HOME=/usr/lib/jvm/java-17-openjdk
ExecStart=${INSTALL_DIR}/kafka/bin/kafka-server-start.sh -daemon ${INSTALL_DIR}/kafka/config/server.properties
ExecStop=${INSTALL_DIR}/kafka/bin/kafka-server-stop.sh
TimeoutSec=30
Restart=on-failure
[Install]
WantedBy=multi-user.target
EOF
systemctl daemon-reload
systemctl enable zookeeper kafka
echo "[7/9] Starting Kafka cluster..."
systemctl start zookeeper
sleep 5
systemctl start kafka
sleep 10
echo "[8/9] Creating Maven project structure..."
PROJECT_DIR="/home/$KAFKA_USER/$PROJECT_NAME"
sudo -u "$KAFKA_USER" mkdir -p "$PROJECT_DIR/src/main/java/com/example"
# Create POM file
sudo -u "$KAFKA_USER" cat > "$PROJECT_DIR/pom.xml" << 'EOF'
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<groupId>com.example</groupId>
<artifactId>kafka-streams-optimization</artifactId>
<version>1.0.0</version>
<properties>
<maven.compiler.source>17</maven.compiler.source>
<maven.compiler.target>17</maven.compiler.target>
<kafka.version>3.6.0</kafka.version>
</properties>
<dependencies>
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-streams</artifactId>
<version>${kafka.version}</version>
</dependency>
<dependency>
<groupId>org.rocksdb</groupId>
<artifactId>rocksdbjni</artifactId>
<version>8.8.1</version>
</dependency>
<dependency>
<groupId>org.slf4j</groupId>
<artifactId>slf4j-simple</artifactId>
<version>1.7.36</version>
</dependency>
</dependencies>
</project>
EOF
# Create RocksDB configuration class
sudo -u "$KAFKA_USER" cat > "$PROJECT_DIR/src/main/java/com/example/OptimizedRocksDBConfig.java" << 'EOF'
package com.example;
import org.apache.kafka.streams.state.RocksDBConfigSetter;
import org.rocksdb.*;
import java.util.Map;
public class OptimizedRocksDBConfig implements RocksDBConfigSetter {
@Override
public void setConfig(String storeName, Options options, Map<String, Object> configs) {
// Optimize block cache for better read performance
BlockBasedTableConfig tableConfig = new BlockBasedTableConfig();
tableConfig.setBlockCacheSize(64 * 1024 * 1024); // 64MB block cache
tableConfig.setBlockSize(16 * 1024); // 16KB block size
tableConfig.setCacheIndexAndFilterBlocks(true);
// Configure bloom filter for faster lookups
tableConfig.setFilterPolicy(new BloomFilter(10, false));
options.setTableFormatConfig(tableConfig);
// Optimize write performance
options.setWriteBufferSize(128 * 1024 * 1024); // 128MB write buffer
options.setMaxWriteBufferNumber(4);
options.setMinWriteBufferNumberToMerge(2);
// Configure compaction for better performance
options.setCompressionType(CompressionType.LZ4_COMPRESSION);
options.setLevelCompactionDynamicLevelBytes(true);
options.setMaxBackgroundCompactions(4);
options.setMaxBackgroundFlushes(2);
}
}
EOF
chown -R "$KAFKA_USER:$KAFKA_USER" "$PROJECT_DIR"
echo "[9/9] Verifying installation..."
if systemctl is-active --quiet kafka && systemctl is-active --quiet zookeeper; then
print_status "ZooKeeper and Kafka services are running"
# Test Kafka connectivity
if sudo -u "$KAFKA_USER" "${INSTALL_DIR}/kafka/bin/kafka-broker-api-versions.sh" --bootstrap-server localhost:9092 >/dev/null 2>&1; then
print_status "Kafka cluster is responding"
else
print_warning "Kafka cluster might not be fully ready yet"
fi
print_status "✅ Installation completed successfully!"
echo
echo "Kafka installation details:"
echo "- Install directory: ${INSTALL_DIR}/kafka"
echo "- Kafka user: $KAFKA_USER"
echo "- Project directory: $PROJECT_DIR"
echo "- Bootstrap server: localhost:9092"
echo
echo "Services can be managed with:"
echo "- systemctl start/stop/restart zookeeper"
echo "- systemctl start/stop/restart kafka"
else
print_error "Services are not running properly"
exit 1
fi
Review the script before running. Execute with: bash install.sh