From 187c17c8d1080b13cbb53aa0667595dfc6cc579f Mon Sep 17 00:00:00 2001 From: Haitao Pan Date: Tue, 7 Jan 2025 14:44:54 +0800 Subject: [PATCH] update Kafka-Alert-Demo --- .../KafkaConsumer-Python-Demo.py | 64 ---------------- Solutions/Kafka-Alert-Demo/KafkaConsumer.py | 74 +++++++++++++++++++ .../KafkaProducer-Python-Demo.py | 26 ------- Solutions/Kafka-Alert-Demo/KafkaProducer.py | 49 ++++++++++++ .../Kafka-Alert-Demo/Setup-Kafka-Cluster.sh | 11 --- .../Kafka-Alert-Demo/Setup-Redis-Cluster.sh | 10 --- .../Kafka-Alert-Demo/cn-k3s-all-in-one.sh | 11 +++ Solutions/Kafka-Alert-Demo/set-node-label.sh | 6 ++ .../setup-kafka-cluster-public.sh | 36 +++++++++ .../Kafka-Alert-Demo/setup-kafka-cluster.sh | 23 ++++++ .../Kafka-Alert-Demo/setup-redis-cluster.sh | 16 ++++ 11 files changed, 215 insertions(+), 111 deletions(-) delete mode 100644 Solutions/Kafka-Alert-Demo/KafkaConsumer-Python-Demo.py create mode 100644 Solutions/Kafka-Alert-Demo/KafkaConsumer.py delete mode 100644 Solutions/Kafka-Alert-Demo/KafkaProducer-Python-Demo.py create mode 100644 Solutions/Kafka-Alert-Demo/KafkaProducer.py delete mode 100644 Solutions/Kafka-Alert-Demo/Setup-Kafka-Cluster.sh delete mode 100644 Solutions/Kafka-Alert-Demo/Setup-Redis-Cluster.sh create mode 100644 Solutions/Kafka-Alert-Demo/cn-k3s-all-in-one.sh create mode 100644 Solutions/Kafka-Alert-Demo/set-node-label.sh create mode 100644 Solutions/Kafka-Alert-Demo/setup-kafka-cluster-public.sh create mode 100644 Solutions/Kafka-Alert-Demo/setup-kafka-cluster.sh create mode 100644 Solutions/Kafka-Alert-Demo/setup-redis-cluster.sh diff --git a/Solutions/Kafka-Alert-Demo/KafkaConsumer-Python-Demo.py b/Solutions/Kafka-Alert-Demo/KafkaConsumer-Python-Demo.py deleted file mode 100644 index bc9d556b..00000000 --- a/Solutions/Kafka-Alert-Demo/KafkaConsumer-Python-Demo.py +++ /dev/null @@ -1,64 +0,0 @@ -from kafka import KafkaConsumer -import redis -import json -import smtplib -from email.mime.text import MIMEText -from email.mime.multipart import MIMEMultipart - -# Kafka 和 Redis 配置 -KAFKA_SERVER = 'localhost:9092' -REDIS_HOST = 'localhost' -REDIS_PORT = 6379 -ALARM_TOPIC = 'alarm_topic' - -# 邮件配置 -SMTP_SERVER = 'smtp.example.com' -SMTP_PORT = 587 -EMAIL_ADDRESS = 'alert@example.com' -EMAIL_PASSWORD = 'your_password' -RECIPIENT_EMAIL = 'recipient@example.com' - -# 初始化 Kafka 消费者和 Redis 客户端 -consumer = KafkaConsumer(ALARM_TOPIC, bootstrap_servers=KAFKA_SERVER) -redis_client = redis.StrictRedis(host=REDIS_HOST, port=REDIS_PORT, db=0) - -# 邮件发送函数 -def send_email(alarm): - msg = MIMEMultipart() - msg['From'] = EMAIL_ADDRESS - msg['To'] = RECIPIENT_EMAIL - msg['Subject'] = f"Alarm Notification - {alarm['level']}" - body = f""" - Alarm ID: {alarm['alarm_id']} - Level: {alarm['level']} - Message: {alarm['message']} - Source: {alarm['source']} - Timestamp: {alarm['timestamp']} - """ - msg.attach(MIMEText(body, 'plain')) - with smtplib.SMTP(SMTP_SERVER, SMTP_PORT) as server: - server.starttls() - server.login(EMAIL_ADDRESS, EMAIL_PASSWORD) - server.sendmail(EMAIL_ADDRESS, RECIPIENT_EMAIL, msg.as_string()) - print(f"Email sent for alarm: {alarm['alarm_id']}") - -# 告警处理函数 -def process_alarm(message, deduplication=True): - alarm_data = json.loads(message.value) - alarm_id = alarm_data['alarm_id'] - timestamp = alarm_data['timestamp'] - key = f"{alarm_id}:{timestamp}" - - # 去重逻辑 - if deduplication: - if not redis_client.exists(key): - redis_client.setex(key, 3600, "1") # 设置1小时过期 - send_email(alarm_data) - else: - print(f"Duplicate alarm discarded: {alarm_id}") - else: - send_email(alarm_data) - -# 消费 Kafka 消息 -for message in consumer: - process_alarm(message, deduplication=True) # 控制去重 diff --git a/Solutions/Kafka-Alert-Demo/KafkaConsumer.py b/Solutions/Kafka-Alert-Demo/KafkaConsumer.py new file mode 100644 index 00000000..3d4c6794 --- /dev/null +++ b/Solutions/Kafka-Alert-Demo/KafkaConsumer.py @@ -0,0 +1,74 @@ +from kafka import KafkaConsumer +import json +import logging +import smtplib +from email.mime.text import MIMEText +from email.mime.multipart import MIMEMultipart + +# 配置日志 +logging.basicConfig(level=logging.INFO) + +# Kafka 配置 +KAFKA_SERVER = '8.130.111.218:9092' +ALARM_TOPIC = 'your_topic_name' +KAFKA_GROUP_ID = 'your_consumer_group' # 消费者组 ID + +# 邮件配置 +SMTP_SERVER = 'smtp.qq.com' +SMTP_PORT = 465 # 465端口支持SSL加密 +EMAIL_ADDRESS = 'manbuzhe2009@qq.com' +EMAIL_PASSWORD = 'xxxxxx' # QQ授权码 +RECIPIENT_EMAIL = '156405189@qq.com' + +# Kafka Consumer 配置 +def create_kafka_consumer(): + return KafkaConsumer( + ALARM_TOPIC, # 订阅的 Kafka topic + bootstrap_servers=KAFKA_SERVER, # Kafka 集群地址 + group_id=KAFKA_GROUP_ID, # 消费者组 ID + value_deserializer=lambda x: json.loads(x.decode('utf-8')), # 反序列化消息 + sasl_mechanism='PLAIN', # SASL 认证机制 + sasl_plain_username='user1', # Kafka 认证用户名 + sasl_plain_password='test', # Kafka 认证密码 + security_protocol='SASL_PLAINTEXT', # 安全协议 + auto_offset_reset='earliest' # 从最早的消息开始消费 + ) + +# 消费 Kafka 消息并发送邮件 +def consume_messages_and_send_email(consumer): + for message in consumer: + logging.info(f"Consumed message: {message.value}") # 输出接收到的消息 + # 将消息发送到邮件 + subject = f"Kafka Alert - New message at offset {message.offset}" + body = f"A new message was received in topic '{ALARM_TOPIC}' at offset {message.offset}. The message is:\n\n{json.dumps(message.value, indent=2)}" + send_email(subject, body) + +# 发送邮件 +def send_email(subject, body): + try: + # 设置邮件内容 + msg = MIMEMultipart() + msg['From'] = EMAIL_ADDRESS + msg['To'] = RECIPIENT_EMAIL + msg['Subject'] = subject + msg.attach(MIMEText(body, 'plain')) + + # 连接 SMTP 服务器并发送邮件 + with smtplib.SMTP_SSL(SMTP_SERVER, SMTP_PORT) as server: + server.login(EMAIL_ADDRESS, EMAIL_PASSWORD) + server.sendmail(EMAIL_ADDRESS, RECIPIENT_EMAIL, msg.as_string()) + + logging.info(f"Email sent to {RECIPIENT_EMAIL}") + except Exception as e: + logging.error(f"Failed to send email: {e}") + +# 主函数 +def main(): + # 创建 Kafka 消费者 + consumer = create_kafka_consumer() + # 开始消费 Kafka 消息并发送邮件 + consume_messages_and_send_email(consumer) + +# 启动应用 +if __name__ == "__main__": + main() diff --git a/Solutions/Kafka-Alert-Demo/KafkaProducer-Python-Demo.py b/Solutions/Kafka-Alert-Demo/KafkaProducer-Python-Demo.py deleted file mode 100644 index 53a70ff8..00000000 --- a/Solutions/Kafka-Alert-Demo/KafkaProducer-Python-Demo.py +++ /dev/null @@ -1,26 +0,0 @@ -from kafka import KafkaProducer -import json -import time -import random - -producer = KafkaProducer( - bootstrap_servers='localhost:9092', - value_serializer=lambda v: json.dumps(v).encode('utf-8') -) - -alarm_levels = ["INFO", "WARNING", "CRITICAL"] - -def generate_alarm(): - return { - "alarm_id": random.randint(1000, 9999), - "timestamp": int(time.time()), - "level": random.choice(alarm_levels), - "message": "System load high", - "source": "Server-01" - } - -while True: - alarm = generate_alarm() - producer.send('alarm_topic', alarm) - print(f"Produced: {alarm}") - time.sleep(5) diff --git a/Solutions/Kafka-Alert-Demo/KafkaProducer.py b/Solutions/Kafka-Alert-Demo/KafkaProducer.py new file mode 100644 index 00000000..657f8486 --- /dev/null +++ b/Solutions/Kafka-Alert-Demo/KafkaProducer.py @@ -0,0 +1,49 @@ +from kafka import KafkaProducer +import json +import logging +import time + +# 配置日志 +logging.basicConfig(level=logging.INFO) + +# Kafka Producer 配置 +producer = KafkaProducer( + bootstrap_servers='8.130.111.218:44204', + value_serializer=lambda v: json.dumps(v).encode('utf-8'), + sasl_mechanism='PLAIN', + sasl_plain_username='user1', + sasl_plain_password='test', + security_protocol='SASL_PLAINTEXT', +) + +# 目标 topic +topic = 'your_topic_name' + +# 循环次数不限制,直到手动停止 +attempt = 0 + +# 模拟持续发送不同消息 +while True: + message = {"key": f"value_{attempt}", "status": "success", "attempt": attempt} + + try: + # 发送消息并等待确认 + future = producer.send(topic, value=message) + + # 等待确认并获取结果 + record_metadata = future.get(timeout=10) + + # 输出消息成功写入的元数据 + logging.info(f"Message sent to topic {record_metadata.topic} partition {record_metadata.partition} with offset {record_metadata.offset}") + + except Exception as e: + logging.error(f"Error sending message: {e}") + + # 增加尝试次数 + attempt += 1 + + # 暂停 1 秒钟,确保每次发送的间隔为 1 秒 + time.sleep(1) # 每次发送后暂停 1 秒钟 + +# 关闭 Kafka 生产者连接(如果手动停止程序时才会关闭) +producer.close() diff --git a/Solutions/Kafka-Alert-Demo/Setup-Kafka-Cluster.sh b/Solutions/Kafka-Alert-Demo/Setup-Kafka-Cluster.sh deleted file mode 100644 index 163c4ff3..00000000 --- a/Solutions/Kafka-Alert-Demo/Setup-Kafka-Cluster.sh +++ /dev/null @@ -1,11 +0,0 @@ -helm repo add bitnami https://charts.bitnami.com/bitnami -helm repo update -kubectl create namespace kafka -helm install kafka bitnami/kafka --namespace kafka \ - --set replicaCount=3 \ - --set persistence.enabled=true \ - --set persistence.size=8Gi \ - --set externalZookeeper.enabled=false \ - --set zookeeper.enabled=true -kubectl get pods --namespace kafka -kubectl get svc --namespace kafka diff --git a/Solutions/Kafka-Alert-Demo/Setup-Redis-Cluster.sh b/Solutions/Kafka-Alert-Demo/Setup-Redis-Cluster.sh deleted file mode 100644 index d263ded0..00000000 --- a/Solutions/Kafka-Alert-Demo/Setup-Redis-Cluster.sh +++ /dev/null @@ -1,10 +0,0 @@ -helm repo add bitnami https://charts.bitnami.com/bitnami -helm repo update -kubectl create namespace redis -helm install redis bitnami/redis-cluster --namespace redis \ - --set cluster.enabled=true \ - --set cluster.nodes=6 \ - --set persistence.enabled=true \ - --set persistence.size=8Gi -kubectl get pods --namespace redis -kubectl get svc --namespace redis diff --git a/Solutions/Kafka-Alert-Demo/cn-k3s-all-in-one.sh b/Solutions/Kafka-Alert-Demo/cn-k3s-all-in-one.sh new file mode 100644 index 00000000..3fc2a769 --- /dev/null +++ b/Solutions/Kafka-Alert-Demo/cn-k3s-all-in-one.sh @@ -0,0 +1,11 @@ +sudo mkdir -pv /opt/rancher/k3s +curl -sfL https://rancher-mirror.rancher.cn/k3s/k3s-install.sh | INSTALL_K3S_MIRROR=cn INSTALL_K3S_SKIP_SELINUX_RPM=true sh -s - \ +--system-default-registry "registry.cn-hangzhou.aliyuncs.com" --data-dir=/opt/rancher/k3s --kube-apiserver-arg service-node-port-range=0-50000 +#curl -sfL https://get.k3s.io | sh -s - --disable=traefik,servicelb \ +# --data-dir=/opt/rancher/k3s \ +# --kube-apiserver-arg service-node-port-range=0-50000 + +sudo mkdir -pv ~/.kube/ +sudo cp /etc/rancher/k3s/k3s.yaml ~/.kube/config + +sudo snap install helm --classic diff --git a/Solutions/Kafka-Alert-Demo/set-node-label.sh b/Solutions/Kafka-Alert-Demo/set-node-label.sh new file mode 100644 index 00000000..6a018e03 --- /dev/null +++ b/Solutions/Kafka-Alert-Demo/set-node-label.sh @@ -0,0 +1,6 @@ +k8s_node=`sudo kubectl get nodes | awk 'NR>1{print $1}'` + +sudo kubectl label node $k8s_node master_controller=enable +sudo kubectl label node $k8s_node tsdb=enable +sudo kubectl label node $k8s_node dfdb=enable +sudo kubectl label node $k8s_node elasticsearch-warm=enable diff --git a/Solutions/Kafka-Alert-Demo/setup-kafka-cluster-public.sh b/Solutions/Kafka-Alert-Demo/setup-kafka-cluster-public.sh new file mode 100644 index 00000000..5a736309 --- /dev/null +++ b/Solutions/Kafka-Alert-Demo/setup-kafka-cluster-public.sh @@ -0,0 +1,36 @@ +helm repo add bitnami https://charts.bitnami.com/bitnami +helm repo update +kubectl create namespace kafka || true +helm upgrade --install kafka bitnami/kafka --namespace kafka \ + --set global.security.allowInsecureImages=true \ + --set image.registry='images.onwalk.net' \ + --set image.repository='public/kafka' \ + --set image.tag='3.9.0-debian-12-r4' \ + --set replicaCount=3 \ + --set sasl.enabledMechanisms="PLAIN" \ + --set sasl.interBrokerMechanism=PLAIN \ + --set sasl.controllerMechanism=PLAIN \ + --set rbac.create=true \ + --set externalAccess.enabled=true \ + --set externalAccess.autoDiscovery.enabled=true \ + --set externalAccess.autoDiscovery.image.registry=images.onwalk.net \ + --set externalAccess.autoDiscovery.image.repository=public/kubectl \ + --set externalAccess.autoDiscovery.image.tag=1.32.0-debian-12-r0 \ + --set externalAccess.broker.service.type=NodePort \ + --set externalAccess.broker.service.externalIPs[0]=8.130.127.232 \ + --set externalAccess.broker.service.externalIPs[1]=8.130.127.232 \ + --set externalAccess.broker.service.externalIPs[2]=8.130.127.232 \ + --set externalAccess.controller.service.type=NodePort \ + --set externalAccess.controller.service.externalIPs[0]=8.130.127.232 \ + --set externalAccess.controller.service.externalIPs[1]=8.130.127.232 \ + --set externalAccess.controller.service.externalIPs[2]=8.130.127.232 \ + --set controller.automountServiceAccountToken=true \ + --set broker.automountServiceAccountToken=true \ + --set sasl.client.users[0]=user1 \ + --set sasl.client.passwords="test" \ + --set persistence.enabled=true \ + --set persistence.size=8Gi \ + --set externalZookeeper.enabled=false \ + --set zookeeper.enabled=false +kubectl get pods --namespace kafka +kubectl get svc --namespace kafka diff --git a/Solutions/Kafka-Alert-Demo/setup-kafka-cluster.sh b/Solutions/Kafka-Alert-Demo/setup-kafka-cluster.sh new file mode 100644 index 00000000..1f8d8d08 --- /dev/null +++ b/Solutions/Kafka-Alert-Demo/setup-kafka-cluster.sh @@ -0,0 +1,23 @@ +helm repo add bitnami https://charts.bitnami.com/bitnami +helm repo update +kubectl create namespace kafka || true +helm upgrade --install kafka bitnami/kafka --namespace kafka \ + --set global.security.allowInsecureImages=true \ + --set global.security.allowInsecureImages=true \ + --set image.registry='images.onwalk.net' \ + --set image.repository='public/kafka' \ + --set image.tag='3.9.0-debian-12-r4' \ + --set replicaCount=1 \ + --set sasl.enabledMechanisms="PLAIN" \ + --set sasl.interBrokerMechanism=PLAIN \ + --set sasl.controllerMechanism=PLAIN \ + --set service.type=NodePort \ + --set service.nodePorts.client="9092" \ + --set sasl.client.users[0]=user1 \ + --set sasl.client.passwords="test" \ + --set persistence.enabled=true \ + --set persistence.size=8Gi \ + --set externalZookeeper.enabled=false \ + --set zookeeper.enabled=false +kubectl get pods --namespace kafka +kubectl get svc --namespace kafka diff --git a/Solutions/Kafka-Alert-Demo/setup-redis-cluster.sh b/Solutions/Kafka-Alert-Demo/setup-redis-cluster.sh new file mode 100644 index 00000000..afbe0863 --- /dev/null +++ b/Solutions/Kafka-Alert-Demo/setup-redis-cluster.sh @@ -0,0 +1,16 @@ +helm repo add bitnami https://charts.bitnami.com/bitnami +helm repo update +kubectl create namespace redis +helm upgrade --install redis bitnami/redis --namespace redis \ + --set global.security.allowInsecureImages=true \ + --set architecture=standalone \ + --set image.registry="images.onwalk.net" \ + --set image.repository="public/redis" \ + --set image.tag="7.4.1-debian-12-r3" \ + --set auth.enabled=false \ + --set cluster.enabled=false \ + --set cluster.nodes=1 \ + --set persistence.enabled=true \ + --set persistence.size=8Gi +kubectl get pods --namespace redis +kubectl get svc --namespace redis