В Ansible, помимо модулей, можно писать свои фильтры — filter_plugins. Это по сути обычные Python-функции, которые можно вызывать в шаблонах и тасках.
Например, у нас есть Kafka-роль, в которой понадобилось добавить возможность менять фактор репликации (replication factor). Операция редкая, но почему бы сразу не сделать через IaC?
Kafka предоставляет утилиту
kafka-reassign-partitions.sh, которая формирует JSON с текущим и предложенным состоянием.Примерная команды:
kafka-reassign-partitions.sh --bootstrap-server <ip>:<port> --command-config admin_connect.cfg --generate --topics-to-move-json-file /tmp/kafka_topics_to_move.json --broker-list "1,2,3"
broker-list это id брокеров, не их IP адреса.
В итоге мы получим что-то вроде:
Current partition replica assignment
{"version":1,"partitions":[{"topic":"test-test-test","partition":0,"replicas":[3],"log_dirs":["any"]},{"topic":"test-test-test","partition":1,"replicas":[1],"log_dirs":["any"]},{"topic":"test-test-test","partition":2,"replicas":[2],"log_dirs":["any"]}]}
Proposed partition reassignment configuration
{"version":1,"partitions":[{"topic":"test-test-test","partition":0,"replicas":[2],"log_dirs":["any"]},{"topic":"test-test-test","partition":1,"replicas":[3],"log_dirs":["any"]},{"topic":"test-test-test","partition":2,"replicas":[1],"log_dirs":["any"]}]}
Отсюда нам полезно взять то, что после Proposed. В ansible это можно сделать так:
- name: Extract JSON from reassignment plan
set_fact:
proposed_json: "{{ reassignment_plan.stdout.split('Proposed partition reassignment configuration')[-1] | trim | from_json }}"
А дальше уже интереснее — логичнее использовать кастомную фильтр-функцию, чтобы перераспределить реплики с нужным RF.
Создадим в корне роли директорию filter_plugins с файлом set_new_rf.py:
from typing import Dict, Any
import itertools
class FilterModule(object):
"""Custom filter to set replication factor in kafka."""
def filters(self):
return {
'set_rf': self.set_rf,
}
def set_rf(self, proposed_json: Dict[str, Any], new_rf: int, kafka_brokers: str) -> Dict[str, Any]:
"""Return copy of proposed_json with new replication factor"""
updated_partitions = []
brokers = [int(x) for x in kafka_brokers.split(',')]
broker_cycle = itertools.cycle(brokers)
for partition in proposed_json.get('partitions', []):
replicas = partition['replicas']
unique_replicas = list(set(replicas))
length = len(unique_replicas)
if new_rf <= length:
new_replicas = sorted(unique_replicas, key=lambda x: brokers.index(x))[:new_rf]
else:
needed = new_rf - length
possible_kafka_brokers = [broker for broker in brokers if broker not in unique_replicas]
new_replicas = unique_replicas + possible_kafka_brokers[:needed]
if new_rf == 1:
new_replicas = [next(broker_cycle)]
updated_partitions.append({
'topic': partition['topic'],
'partition': partition['partition'],
'replicas': new_replicas
})
return {
'version': proposed_json.get('version', 1),
'partitions': updated_partitions
}
Код сформирует json с равномерным распределением партиций как при увеличении RF, так и при снижении.
Вызвать в таске можно например так:
- name: Modify proposed json via custom filter
set_fact:
updated_proposed_json: "{{ proposed_json | set_rf(NEW_RF|int, kafka_broker_list) }}"
Остается только выполнить execute:
kafka-reassign-partitions.sh --bootstrap-server <ip>:<port> --execute --reassignment-json-file /tmp/kafka_updated_proposed.json --command-config admin_connect.cfg
#ansible