update Kafka-Alert-Demo
This commit is contained in:
parent
a911377a1f
commit
187c17c8d1
@ -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) # 控制去重
|
||||
74
Solutions/Kafka-Alert-Demo/KafkaConsumer.py
Normal file
74
Solutions/Kafka-Alert-Demo/KafkaConsumer.py
Normal file
@ -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()
|
||||
@ -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)
|
||||
49
Solutions/Kafka-Alert-Demo/KafkaProducer.py
Normal file
49
Solutions/Kafka-Alert-Demo/KafkaProducer.py
Normal file
@ -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()
|
||||
@ -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
|
||||
@ -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
|
||||
11
Solutions/Kafka-Alert-Demo/cn-k3s-all-in-one.sh
Normal file
11
Solutions/Kafka-Alert-Demo/cn-k3s-all-in-one.sh
Normal file
@ -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
|
||||
6
Solutions/Kafka-Alert-Demo/set-node-label.sh
Normal file
6
Solutions/Kafka-Alert-Demo/set-node-label.sh
Normal file
@ -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
|
||||
36
Solutions/Kafka-Alert-Demo/setup-kafka-cluster-public.sh
Normal file
36
Solutions/Kafka-Alert-Demo/setup-kafka-cluster-public.sh
Normal file
@ -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
|
||||
23
Solutions/Kafka-Alert-Demo/setup-kafka-cluster.sh
Normal file
23
Solutions/Kafka-Alert-Demo/setup-kafka-cluster.sh
Normal file
@ -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
|
||||
16
Solutions/Kafka-Alert-Demo/setup-redis-cluster.sh
Normal file
16
Solutions/Kafka-Alert-Demo/setup-redis-cluster.sh
Normal file
@ -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
|
||||
Loading…
Reference in New Issue
Block a user