From f375304e3f06dbd6f2aefe2563fa8b301805bdc0 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E2=80=9CNaeel=E2=80=9D?= Date: Tue, 21 Jul 2026 17:42:42 +0400 Subject: [PATCH] fix: SASL_SSL + TLS certs for Kafka --- app.py | 61 +++++++++++++++++++++++++++++++++++++++++++--------------- 1 file changed, 45 insertions(+), 16 deletions(-) diff --git a/app.py b/app.py index 8a2fc02..d58f346 100644 --- a/app.py +++ b/app.py @@ -1,12 +1,15 @@ """ -IoT Consumer — Flask app that reads IoT events from Kafka +IoT Consumer — Flask app that reads IoT events from Kafka (TLS + SASL) and inserts them into ClickHouse. Env vars (set by Terraform): - KAFKA_BROKERS — Kafka bootstrap servers + KAFKA_BROKERS — bootstrap server (host:port) KAFKA_TOPIC — topic name KAFKA_USERNAME — SASL username KAFKA_PASSWORD — SASL password + KAFKA_CA_CRT — CA certificate (PEM) + KAFKA_USER_CRT — user certificate (PEM) + KAFKA_USER_KEY — user private key (PEM) KAFKA_GROUP_ID — consumer group id CH_HOST — ClickHouse host CH_PORT — ClickHouse port (default 8123) @@ -17,6 +20,7 @@ Env vars (set by Terraform): import json import os +import tempfile import threading from flask import Flask, jsonify @@ -26,10 +30,13 @@ import clickhouse_connect app = Flask(__name__) # --- Config --- -KAFKA_BROKERS = os.environ["KAFKA_BROKERS"].split(",") +KAFKA_BROKERS = os.environ["KAFKA_BROKERS"] KAFKA_TOPIC = os.environ["KAFKA_TOPIC"] KAFKA_USERNAME = os.environ["KAFKA_USERNAME"] KAFKA_PASSWORD = os.environ["KAFKA_PASSWORD"] +KAFKA_CA_CRT = os.environ.get("KAFKA_CA_CRT", "") +KAFKA_USER_CRT = os.environ.get("KAFKA_USER_CRT", "") +KAFKA_USER_KEY = os.environ.get("KAFKA_USER_KEY", "") KAFKA_GROUP_ID = os.environ.get("KAFKA_GROUP_ID", "iot-consumer-group") CH_HOST = os.environ["CH_HOST"] @@ -38,6 +45,21 @@ CH_USER = os.environ["CH_USER"] CH_PASSWORD = os.environ["CH_PASSWORD"] CH_DATABASE = os.environ["CH_DATABASE"] +# --- TLS certs → temp files --- +_cert_dir = tempfile.mkdtemp(prefix="kafka-certs-") + +def _write_cert(name, content): + if not content: + return None + path = os.path.join(_cert_dir, name) + with open(path, "w") as f: + f.write(content) + return path + +_ca_file = _write_cert("ca.crt", KAFKA_CA_CRT) +_cert_file = _write_cert("user.crt", KAFKA_USER_CRT) +_key_file = _write_cert("user.key", KAFKA_USER_KEY) + events_consumed = 0 errors_count = 0 ch_client = None @@ -53,7 +75,6 @@ def get_ch_client(): password=CH_PASSWORD, database=CH_DATABASE, ) - # Create table if not exists ch_client.command(""" CREATE TABLE IF NOT EXISTS iot_events ( event_id String, @@ -94,18 +115,26 @@ def insert_batch(events): def consume_loop(): global events_consumed, errors_count - consumer = KafkaConsumer( - KAFKA_TOPIC, - bootstrap_servers=KAFKA_BROKERS, - security_protocol="SASL_PLAINTEXT", - sasl_mechanism="PLAIN", - sasl_plain_username=KAFKA_USERNAME, - sasl_plain_password=KAFKA_PASSWORD, - group_id=KAFKA_GROUP_ID, - auto_offset_reset="earliest", - value_deserializer=lambda m: json.loads(m.decode("utf-8")), - max_poll_records=20, - ) + kwargs = { + "bootstrap_servers": KAFKA_BROKERS, + "security_protocol": "SASL_SSL", + "sasl_mechanism": "PLAIN", + "sasl_plain_username": KAFKA_USERNAME, + "sasl_plain_password": KAFKA_PASSWORD, + "group_id": KAFKA_GROUP_ID, + "auto_offset_reset": "earliest", + "value_deserializer": lambda m: json.loads(m.decode("utf-8")), + "max_poll_records": 20, + } + if _ca_file: + kwargs["ssl_cafile"] = _ca_file + if _cert_file: + kwargs["ssl_certfile"] = _cert_file + if _key_file: + kwargs["ssl_keyfile"] = _key_file + kwargs["ssl_check_hostname"] = False + + consumer = KafkaConsumer(KAFKA_TOPIC, **kwargs) print(f"[CONSUMER] Listening on topic '{KAFKA_TOPIC}'...")