Introduction
Apache Kafka is the de facto standard for distributed event streaming — handling real-time data pipelines, event sourcing, and log aggregation at scale. Ansible automates the full deployment: ZooKeeper or KRaft mode clusters, broker configuration, topic management, TLS/SASL security, and monitoring via JMX and Prometheus.
Deploy Kafka with KRaft (No ZooKeeper)
---
- name: Deploy Kafka KRaft cluster
hosts: kafka_brokers
become: true
vars:
kafka_version: "3.8.0"
scala_version: "2.13"
kafka_home: /opt/kafka
kafka_data_dir: /data/kafka
kafka_cluster_id: "{{ lookup('pipe', 'cat /proc/sys/kernel/random/uuid') | regex_replace('-', '') | truncate(22, true, '') }}"
kafka_heap_size: "-Xmx2G -Xms2G"
tasks:
- name: Install Java
ansible.builtin.package:
name: openjdk-17-jre-headless
state: present
- name: Create Kafka user
ansible.builtin.user:
name: kafka
system: true
shell: /usr/sbin/nologin
home: "{{ kafka_home }}"
- name: Download Kafka
ansible.builtin.get_url:
url: "https://downloads.apache.org/kafka/{{ kafka_version }}/kafka_{{ scala_version }}-{{ kafka_version }}.tgz"
dest: /tmp/kafka.tgz
- name: Extract Kafka
ansible.builtin.unarchive:
src: /tmp/kafka.tgz
dest: /opt/
remote_src: true
creates: "{{ kafka_home }}"
- name: Create symlink
ansible.builtin.file:
src: "/opt/kafka_{{ scala_version }}-{{ kafka_version }}"
dest: "{{ kafka_home }}"
state: link
- name: Create data directory
ansible.builtin.file:
path: "{{ kafka_data_dir }}"
state: directory
owner: kafka
mode: '0750'
- name: Deploy KRaft config
ansible.builtin.template:
src: kraft-server.properties.j2
dest: "{{ kafka_home }}/config/kraft/server.properties"
owner: kafka
mode: '0640'
notify: restart kafka
- name: Format storage (first time only)
ansible.builtin.command: >
{{ kafka_home }}/bin/kafka-storage.sh format
-t {{ kafka_cluster_id }}
-c {{ kafka_home }}/config/kraft/server.properties
args:
creates: "{{ kafka_data_dir }}/meta.properties"
become_user: kafka
- name: Create systemd service
ansible.builtin.copy:
dest: /etc/systemd/system/kafka.service
content: |
[Unit]
Description=Apache Kafka
After=network-online.target
[Service]
User=kafka
Environment="KAFKA_HEAP_OPTS={{ kafka_heap_size }}"
ExecStart={{ kafka_home }}/bin/kafka-server-start.sh {{ kafka_home }}/config/kraft/server.properties
ExecStop={{ kafka_home }}/bin/kafka-server-stop.sh
Restart=on-failure
RestartSec=10
LimitNOFILE=65536
[Install]
WantedBy=multi-user.target
mode: '0644'
notify:
- daemon reload
- restart kafka
- name: Allow Kafka through firewall
ansible.posix.firewalld:
port: "{{ item }}/tcp"
permanent: true
state: enabled
immediate: true
loop: ["9092", "9093"]
- name: Start Kafka
ansible.builtin.service:
name: kafka
state: started
enabled: true
handlers:
- name: daemon reload
ansible.builtin.systemd:
daemon_reload: true
- name: restart kafka
ansible.builtin.service:
name: kafka
state: restarted
KRaft Config Template
# templates/kraft-server.properties.j2
# KRaft mode — no ZooKeeper
process.roles=broker,controller
node.id={{ groups['kafka_brokers'].index(inventory_hostname) + 1 }}
controller.quorum.voters={% for host in groups['kafka_brokers'] %}{{ loop.index }}@{{ hostvars[host].ansible_default_ipv4.address }}:9093{% if not loop.last %},{% endif %}{% endfor %}
listeners=PLAINTEXT://:9092,CONTROLLER://:9093
inter.broker.listener.name=PLAINTEXT
controller.listener.names=CONTROLLER
advertised.listeners=PLAINTEXT://{{ ansible_default_ipv4.address }}:9092
log.dirs={{ kafka_data_dir }}
num.partitions=3
default.replication.factor={{ [groups['kafka_brokers'] | length, 3] | min }}
min.insync.replicas={{ [groups['kafka_brokers'] | length - 1, 2] | min }}
log.retention.hours=168
log.segment.bytes=1073741824
log.retention.check.interval.ms=300000
num.network.threads=8
num.io.threads=16
socket.send.buffer.bytes=102400
socket.receive.buffer.bytes=102400
socket.request.max.bytes=104857600
offsets.topic.replication.factor={{ [groups['kafka_brokers'] | length, 3] | min }}
transaction.state.log.replication.factor={{ [groups['kafka_brokers'] | length, 3] | min }}
transaction.state.log.min.isr={{ [groups['kafka_brokers'] | length - 1, 2] | min }}
auto.create.topics.enable=false
delete.topic.enable=true
Manage Topics
- name: Create Kafka topics
ansible.builtin.command: >
{{ kafka_home }}/bin/kafka-topics.sh
--bootstrap-server localhost:9092
--create
--topic {{ item.name }}
--partitions {{ item.partitions | default(6) }}
--replication-factor {{ item.replication | default(3) }}
--config retention.ms={{ item.retention_ms | default(604800000) }}
loop:
- { name: events, partitions: 12, replication: 3 }
- { name: logs, partitions: 6, replication: 3, retention_ms: 259200000 }
- { name: metrics, partitions: 6, replication: 3, retention_ms: 86400000 }
- { name: dead-letter, partitions: 3, replication: 3 }
register: topic_result
changed_when: "'Created topic' in topic_result.stdout"
failed_when: false
run_once: true
- name: List topics
ansible.builtin.command: >
{{ kafka_home }}/bin/kafka-topics.sh
--bootstrap-server localhost:9092 --list
register: topic_list
changed_when: false
run_once: true
- name: Alter topic config
ansible.builtin.command: >
{{ kafka_home }}/bin/kafka-configs.sh
--bootstrap-server localhost:9092
--alter --entity-type topics --entity-name {{ item.name }}
--add-config retention.ms={{ item.retention_ms }}
loop:
- { name: logs, retention_ms: 172800000 }
run_once: true
changed_when: true
TLS/SASL Security
- name: Deploy Kafka TLS keystore
ansible.builtin.copy:
src: "{{ kafka_keystore_file }}"
dest: "{{ kafka_home }}/config/kafka.keystore.jks"
owner: kafka
mode: '0600'
notify: restart kafka
- name: Deploy Kafka TLS truststore
ansible.builtin.copy:
src: "{{ kafka_truststore_file }}"
dest: "{{ kafka_home }}/config/kafka.truststore.jks"
owner: kafka
mode: '0600'
notify: restart kafka
Add to server.properties:
# TLS configuration
listeners=SSL://:9093,CONTROLLER://:9094
ssl.keystore.location={{ kafka_home }}/config/kafka.keystore.jks
ssl.keystore.password={{ vault_kafka_keystore_password }}
ssl.key.password={{ vault_kafka_key_password }}
ssl.truststore.location={{ kafka_home }}/config/kafka.truststore.jks
ssl.truststore.password={{ vault_kafka_truststore_password }}
ssl.client.auth=required
security.inter.broker.protocol=SSL
Monitoring with JMX
- name: Enable JMX
ansible.builtin.lineinfile:
path: /etc/systemd/system/kafka.service
regexp: '^Environment="KAFKA_JMX'
line: 'Environment="KAFKA_JMX_OPTS=-Dcom.sun.management.jmxremote -Dcom.sun.management.jmxremote.port=9999 -Dcom.sun.management.jmxremote.authenticate=false -Dcom.sun.management.jmxremote.ssl=false"'
insertafter: 'KAFKA_HEAP_OPTS'
notify:
- daemon reload
- restart kafka
Health Check
- name: Check broker status
ansible.builtin.command: >
{{ kafka_home }}/bin/kafka-metadata.sh
--snapshot {{ kafka_data_dir }}/__cluster_metadata-0/00000000000000000000.log
--cluster-id {{ kafka_cluster_id }}
register: broker_status
changed_when: false
run_once: true
- name: Test produce/consume
block:
- name: Produce test message
ansible.builtin.shell: >
echo "test-{{ ansible_date_time.epoch }}" |
{{ kafka_home }}/bin/kafka-console-producer.sh
--bootstrap-server localhost:9092 --topic events
changed_when: true
- name: Consume test message
ansible.builtin.command: >
{{ kafka_home }}/bin/kafka-console-consumer.sh
--bootstrap-server localhost:9092 --topic events
--from-beginning --max-messages 1 --timeout-ms 10000
register: consume_result
changed_when: false
run_once: true
Troubleshooting
Under-Replicated Partitions
- name: Check under-replicated partitions
ansible.builtin.command: >
{{ kafka_home }}/bin/kafka-topics.sh
--bootstrap-server localhost:9092
--describe --under-replicated-partitions
register: urp
changed_when: false
run_once: true
- name: Alert on under-replicated partitions
ansible.builtin.debug:
msg: "⚠️ Under-replicated partitions detected!"
when: urp.stdout | length > 0
Consumer Lag
- name: Check consumer group lag
ansible.builtin.command: >
{{ kafka_home }}/bin/kafka-consumer-groups.sh
--bootstrap-server localhost:9092
--describe --group {{ consumer_group }}
register: consumer_lag
changed_when: false
run_once: true
Related Articles
Conclusion
Ansible deploys Kafka clusters in KRaft mode (no ZooKeeper dependency), manages topics with partitions and replication, configures TLS/SASL security, and enables JMX monitoring. Template broker configs from inventory — adding a node means adding a host to the group. Use Kafka for event streaming, log aggregation, and real-time data pipelines, all version-controlled as code.