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

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.