From 0c36d217b8f342ecb3af2dc8ad852dbde655a332 Mon Sep 17 00:00:00 2001 From: Haitao Pan Date: Sat, 11 Jan 2025 08:42:27 +0000 Subject: [PATCH] Tested: KafkaConsumer.py KafkaProducer.py --- Solutions/Kafka-Alert-Demo/KafkaConsumer.py | 30 ++++++++++++++----- Solutions/Kafka-Alert-Demo/KafkaProducer.py | 2 +- Solutions/Kafka-Alert-Demo/Readme.md | 11 +++++++ .../setup-kafka-cluster-private.sh | 12 ++------ .../setup-kafka-cluster-public.sh | 19 +++++++----- 5 files changed, 49 insertions(+), 25 deletions(-) diff --git a/Solutions/Kafka-Alert-Demo/KafkaConsumer.py b/Solutions/Kafka-Alert-Demo/KafkaConsumer.py index 3d4c6794..90a7aeef 100644 --- a/Solutions/Kafka-Alert-Demo/KafkaConsumer.py +++ b/Solutions/Kafka-Alert-Demo/KafkaConsumer.py @@ -9,7 +9,7 @@ from email.mime.multipart import MIMEMultipart logging.basicConfig(level=logging.INFO) # Kafka 配置 -KAFKA_SERVER = '8.130.111.218:9092' +KAFKA_SERVER = '10.43.16.127:9092' ALARM_TOPIC = 'your_topic_name' KAFKA_GROUP_ID = 'your_consumer_group' # 消费者组 ID @@ -17,7 +17,7 @@ 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授权码 +EMAIL_PASSWORD = 'xxxxxxxxxxxxxxxxxx' # QQ授权码 RECIPIENT_EMAIL = '156405189@qq.com' # Kafka Consumer 配置 @@ -36,12 +36,28 @@ def create_kafka_consumer(): # 消费 Kafka 消息并发送邮件 def consume_messages_and_send_email(consumer): + """ + 消费 Kafka 消息并发送邮件,每接收到一条新消息,就发送邮件。 + """ 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) + try: + # 打印日志,显示接收到的消息 + 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}' " + f"at offset {message.offset}.\n\n" + f"Message Content:\n{json.dumps(message.value, indent=2)}" + ) + + # 发送邮件 + send_email(subject, body) + except Exception as e: + logging.error(f"Error processing message: {e}") + + # 发送邮件 def send_email(subject, body): diff --git a/Solutions/Kafka-Alert-Demo/KafkaProducer.py b/Solutions/Kafka-Alert-Demo/KafkaProducer.py index 657f8486..6c097796 100644 --- a/Solutions/Kafka-Alert-Demo/KafkaProducer.py +++ b/Solutions/Kafka-Alert-Demo/KafkaProducer.py @@ -8,7 +8,7 @@ logging.basicConfig(level=logging.INFO) # Kafka Producer 配置 producer = KafkaProducer( - bootstrap_servers='8.130.111.218:44204', + bootstrap_servers='10.43.16.127:9092', value_serializer=lambda v: json.dumps(v).encode('utf-8'), sasl_mechanism='PLAIN', sasl_plain_username='user1', diff --git a/Solutions/Kafka-Alert-Demo/Readme.md b/Solutions/Kafka-Alert-Demo/Readme.md index 29cb964b..71b1bb9e 100644 --- a/Solutions/Kafka-Alert-Demo/Readme.md +++ b/Solutions/Kafka-Alert-Demo/Readme.md @@ -14,3 +14,14 @@ - Kafka-Python 库 - Redis (可选,用于去重) - smtplib (Python 标准库) + +# Install + +apt install python3-pip python3.12-venv -y + +# Create a virtual environment +python3 -m venv kafka-env + +# Activate the virtual environment +source kafka-env/bin/activate +pip install kafka-python-ng diff --git a/Solutions/Kafka-Alert-Demo/setup-kafka-cluster-private.sh b/Solutions/Kafka-Alert-Demo/setup-kafka-cluster-private.sh index c02b0ae5..36190471 100644 --- a/Solutions/Kafka-Alert-Demo/setup-kafka-cluster-private.sh +++ b/Solutions/Kafka-Alert-Demo/setup-kafka-cluster-private.sh @@ -1,7 +1,8 @@ helm repo add bitnami https://charts.bitnami.com/bitnami -helm repo update +#helm repo update kubectl create namespace kafka || true -helm upgrade --install kafka bitnami/kafka --namespace kafka \ +#helm upgrade --install kafka bitnami/kafka --namespace kafka \ +helm upgrade --install kafka kafka-31.1.1.tgz --namespace kafka \ --set global.security.allowInsecureImages=true \ --set image.registry='images.onwalk.net' \ --set image.repository='public/kafka' \ @@ -11,13 +12,6 @@ helm upgrade --install kafka bitnami/kafka --namespace kafka \ --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 controller.automountServiceAccountToken=true \ - --set broker.automountServiceAccountToken=true \ --set sasl.client.users[0]=user1 \ --set sasl.client.passwords="test" \ --set persistence.enabled=true \ diff --git a/Solutions/Kafka-Alert-Demo/setup-kafka-cluster-public.sh b/Solutions/Kafka-Alert-Demo/setup-kafka-cluster-public.sh index 5a736309..6f734cf2 100644 --- a/Solutions/Kafka-Alert-Demo/setup-kafka-cluster-public.sh +++ b/Solutions/Kafka-Alert-Demo/setup-kafka-cluster-public.sh @@ -1,7 +1,8 @@ helm repo add bitnami https://charts.bitnami.com/bitnami -helm repo update +#helm repo update kubectl create namespace kafka || true -helm upgrade --install kafka bitnami/kafka --namespace kafka \ +#helm upgrade --install kafka bitnami/kafka --namespace kafka \ +helm upgrade --install kafka kafka-31.1.1.tgz --namespace kafka \ --set global.security.allowInsecureImages=true \ --set image.registry='images.onwalk.net' \ --set image.repository='public/kafka' \ @@ -11,19 +12,21 @@ helm upgrade --install kafka bitnami/kafka --namespace kafka \ --set sasl.interBrokerMechanism=PLAIN \ --set sasl.controllerMechanism=PLAIN \ --set rbac.create=true \ + --set service.type=NodePort \ + --set service.nodePorts.client="9092" \ --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.broker.service.externalIPs[0]=192.168.80.128 \ + --set externalAccess.broker.service.externalIPs[1]=192.168.80.128 \ + --set externalAccess.broker.service.externalIPs[2]=192.168.80.128 \ --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 externalAccess.controller.service.externalIPs[0]=192.168.80.128 \ + --set externalAccess.controller.service.externalIPs[1]=192.168.80.128 \ + --set externalAccess.controller.service.externalIPs[2]=192.168.80.128 \ --set controller.automountServiceAccountToken=true \ --set broker.automountServiceAccountToken=true \ --set sasl.client.users[0]=user1 \