fix: Kafka → Redis + MongoDB (no ClickHouse)
This commit is contained in:
+2
-1
@@ -1,4 +1,5 @@
|
|||||||
Flask>=3.0
|
Flask>=3.0
|
||||||
gunicorn>=21.2
|
gunicorn>=21.2
|
||||||
kafka-python>=2.0
|
kafka-python>=2.0
|
||||||
clickhouse-connect>=0.7
|
redis>=5.0
|
||||||
|
pymongo>=4.0
|
||||||
|
|||||||
+72
-72
@@ -1,6 +1,6 @@
|
|||||||
"""
|
"""
|
||||||
IoT Consumer — Flask app that reads IoT events from Kafka (TLS + SASL)
|
IoT Consumer — Flask app that reads IoT events from Kafka (TLS + SASL)
|
||||||
and inserts them into ClickHouse.
|
and stores them in Redis (cache) + MongoDB (permanent).
|
||||||
|
|
||||||
Env vars (set by Terraform):
|
Env vars (set by Terraform):
|
||||||
KAFKA_BROKERS — bootstrap server (host:port)
|
KAFKA_BROKERS — bootstrap server (host:port)
|
||||||
@@ -11,11 +11,11 @@ Env vars (set by Terraform):
|
|||||||
KAFKA_USER_CRT — user certificate (PEM)
|
KAFKA_USER_CRT — user certificate (PEM)
|
||||||
KAFKA_USER_KEY — user private key (PEM)
|
KAFKA_USER_KEY — user private key (PEM)
|
||||||
KAFKA_GROUP_ID — consumer group id
|
KAFKA_GROUP_ID — consumer group id
|
||||||
CH_HOST — ClickHouse host
|
REDIS_HOST — Redis host
|
||||||
CH_PORT — ClickHouse port (default 8123)
|
REDIS_PORT — Redis port (default 6379)
|
||||||
CH_USER — ClickHouse user
|
REDIS_PASSWORD — Redis password
|
||||||
CH_PASSWORD — ClickHouse password
|
MONGO_URI — MongoDB connection URI
|
||||||
CH_DATABASE — ClickHouse database name
|
MONGO_DB — MongoDB database name
|
||||||
"""
|
"""
|
||||||
|
|
||||||
import json
|
import json
|
||||||
@@ -25,7 +25,8 @@ import threading
|
|||||||
|
|
||||||
from flask import Flask, jsonify
|
from flask import Flask, jsonify
|
||||||
from kafka import KafkaConsumer
|
from kafka import KafkaConsumer
|
||||||
import clickhouse_connect
|
import redis
|
||||||
|
import pymongo
|
||||||
|
|
||||||
app = Flask(__name__)
|
app = Flask(__name__)
|
||||||
|
|
||||||
@@ -39,11 +40,12 @@ KAFKA_USER_CRT = os.environ.get("KAFKA_USER_CRT", "")
|
|||||||
KAFKA_USER_KEY = os.environ.get("KAFKA_USER_KEY", "")
|
KAFKA_USER_KEY = os.environ.get("KAFKA_USER_KEY", "")
|
||||||
KAFKA_GROUP_ID = os.environ.get("KAFKA_GROUP_ID", "iot-consumer-group")
|
KAFKA_GROUP_ID = os.environ.get("KAFKA_GROUP_ID", "iot-consumer-group")
|
||||||
|
|
||||||
CH_HOST = os.environ["CH_HOST"]
|
REDIS_HOST = os.environ["REDIS_HOST"]
|
||||||
CH_PORT = int(os.environ.get("CH_PORT", 8123))
|
REDIS_PORT = int(os.environ.get("REDIS_PORT", 6379))
|
||||||
CH_USER = os.environ["CH_USER"]
|
REDIS_PASSWORD = os.environ.get("REDIS_PASSWORD", "")
|
||||||
CH_PASSWORD = os.environ["CH_PASSWORD"]
|
|
||||||
CH_DATABASE = os.environ["CH_DATABASE"]
|
MONGO_URI = os.environ["MONGO_URI"]
|
||||||
|
MONGO_DB = os.environ.get("MONGO_DB", "iot")
|
||||||
|
|
||||||
# --- TLS certs → temp files ---
|
# --- TLS certs → temp files ---
|
||||||
_cert_dir = tempfile.mkdtemp(prefix="kafka-certs-")
|
_cert_dir = tempfile.mkdtemp(prefix="kafka-certs-")
|
||||||
@@ -60,56 +62,68 @@ _ca_file = _write_cert("ca.crt", KAFKA_CA_CRT)
|
|||||||
_cert_file = _write_cert("user.crt", KAFKA_USER_CRT)
|
_cert_file = _write_cert("user.crt", KAFKA_USER_CRT)
|
||||||
_key_file = _write_cert("user.key", KAFKA_USER_KEY)
|
_key_file = _write_cert("user.key", KAFKA_USER_KEY)
|
||||||
|
|
||||||
|
# --- Clients ---
|
||||||
events_consumed = 0
|
events_consumed = 0
|
||||||
errors_count = 0
|
errors_count = 0
|
||||||
ch_client = None
|
redis_client = None
|
||||||
|
mongo_client = None
|
||||||
|
mongo_coll = None
|
||||||
|
|
||||||
|
|
||||||
def get_ch_client():
|
def get_redis():
|
||||||
global ch_client
|
global redis_client
|
||||||
if ch_client is None:
|
if redis_client is None:
|
||||||
ch_client = clickhouse_connect.get_client(
|
redis_client = redis.Redis(
|
||||||
host=CH_HOST,
|
host=REDIS_HOST,
|
||||||
port=CH_PORT,
|
port=REDIS_PORT,
|
||||||
username=CH_USER,
|
password=REDIS_PASSWORD or None,
|
||||||
password=CH_PASSWORD,
|
decode_responses=True,
|
||||||
database=CH_DATABASE,
|
|
||||||
)
|
)
|
||||||
ch_client.command("""
|
print("[CONSUMER] Redis connected")
|
||||||
CREATE TABLE IF NOT EXISTS iot_events (
|
return redis_client
|
||||||
event_id String,
|
|
||||||
sensor_id String,
|
|
||||||
device_type String,
|
|
||||||
location String,
|
|
||||||
value Float64,
|
|
||||||
unit String,
|
|
||||||
timestamp DateTime64(3)
|
|
||||||
) ENGINE = MergeTree()
|
|
||||||
ORDER BY (timestamp, sensor_id)
|
|
||||||
""")
|
|
||||||
print("[CONSUMER] ClickHouse table ready")
|
|
||||||
return ch_client
|
|
||||||
|
|
||||||
|
|
||||||
def insert_batch(events):
|
def get_mongo():
|
||||||
if not events:
|
global mongo_client, mongo_coll
|
||||||
return
|
if mongo_client is None:
|
||||||
rows = [[
|
mongo_client = pymongo.MongoClient(MONGO_URI)
|
||||||
e["event_id"],
|
db = mongo_client[MONGO_DB]
|
||||||
e["sensor_id"],
|
mongo_coll = db["iot_events"]
|
||||||
e["device_type"],
|
mongo_coll.create_index("timestamp", expireAfterSeconds=604800)
|
||||||
e["location"],
|
print("[CONSUMER] MongoDB connected")
|
||||||
e["value"],
|
return mongo_coll
|
||||||
e["unit"],
|
|
||||||
e["timestamp"],
|
|
||||||
] for e in events]
|
def store_event(event):
|
||||||
try:
|
r = get_redis()
|
||||||
client = get_ch_client()
|
sensor_id = event["sensor_id"]
|
||||||
client.insert("iot_events", rows,
|
value = event["value"]
|
||||||
column_names=["event_id", "sensor_id", "device_type", "location", "value", "unit", "timestamp"])
|
unit = event["unit"]
|
||||||
except Exception as e:
|
ts = event["timestamp"]
|
||||||
print(f"[CONSUMER] Insert error: {e}")
|
|
||||||
raise
|
# Redis: latest value per sensor
|
||||||
|
r.hset("iot:latest", sensor_id, json.dumps({
|
||||||
|
"value": value, "unit": unit, "location": event["location"], "timestamp": ts,
|
||||||
|
}))
|
||||||
|
|
||||||
|
# Redis: counter per sensor type
|
||||||
|
r.hincrby("iot:counters", event["device_type"], 1)
|
||||||
|
|
||||||
|
# Redis: sorted set for recent events (last 1000)
|
||||||
|
r.zadd("iot:recent", {json.dumps(event): 0})
|
||||||
|
r.zremrangebyrank("iot:recent", 0, -1001)
|
||||||
|
|
||||||
|
# MongoDB: permanent storage
|
||||||
|
coll = get_mongo()
|
||||||
|
coll.insert_one({
|
||||||
|
"event_id": event["event_id"],
|
||||||
|
"sensor_id": sensor_id,
|
||||||
|
"device_type": event["device_type"],
|
||||||
|
"location": event["location"],
|
||||||
|
"value": value,
|
||||||
|
"unit": unit,
|
||||||
|
"timestamp": ts,
|
||||||
|
})
|
||||||
|
|
||||||
|
|
||||||
def consume_loop():
|
def consume_loop():
|
||||||
@@ -135,20 +149,14 @@ def consume_loop():
|
|||||||
kwargs["ssl_check_hostname"] = False
|
kwargs["ssl_check_hostname"] = False
|
||||||
|
|
||||||
consumer = KafkaConsumer(KAFKA_TOPIC, **kwargs)
|
consumer = KafkaConsumer(KAFKA_TOPIC, **kwargs)
|
||||||
|
|
||||||
print(f"[CONSUMER] Listening on topic '{KAFKA_TOPIC}'...")
|
print(f"[CONSUMER] Listening on topic '{KAFKA_TOPIC}'...")
|
||||||
|
|
||||||
batch = []
|
|
||||||
for message in consumer:
|
for message in consumer:
|
||||||
try:
|
try:
|
||||||
event = message.value
|
store_event(message.value)
|
||||||
batch.append(event)
|
|
||||||
events_consumed += 1
|
events_consumed += 1
|
||||||
|
if events_consumed % 10 == 0:
|
||||||
if len(batch) >= 10:
|
print(f"[CONSUMER] Processed {events_consumed} events")
|
||||||
insert_batch(batch)
|
|
||||||
print(f"[CONSUMER] Inserted {len(batch)} events into ClickHouse")
|
|
||||||
batch = []
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
errors_count += 1
|
errors_count += 1
|
||||||
print(f"[CONSUMER] Error: {e}")
|
print(f"[CONSUMER] Error: {e}")
|
||||||
@@ -156,19 +164,11 @@ def consume_loop():
|
|||||||
|
|
||||||
@app.route("/")
|
@app.route("/")
|
||||||
def index():
|
def index():
|
||||||
query = "SELECT count() FROM iot_events"
|
|
||||||
try:
|
|
||||||
client = get_ch_client()
|
|
||||||
total = client.query(query).result_rows[0][0]
|
|
||||||
except Exception:
|
|
||||||
total = 0
|
|
||||||
|
|
||||||
return jsonify({
|
return jsonify({
|
||||||
"service": "iot-consumer",
|
"service": "iot-consumer",
|
||||||
"status": "running",
|
"status": "running",
|
||||||
"kafka_topic": KAFKA_TOPIC,
|
"kafka_topic": KAFKA_TOPIC,
|
||||||
"events_consumed": events_consumed,
|
"events_consumed": events_consumed,
|
||||||
"events_in_db": total,
|
|
||||||
"errors": errors_count,
|
"errors": errors_count,
|
||||||
})
|
})
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user