Tested: KafkaConsumer.py KafkaProducer.py

This commit is contained in:
Haitao Pan 2025-01-11 08:42:27 +00:00
parent 2a86d5683b
commit 0c36d217b8
5 changed files with 49 additions and 25 deletions

View File

@ -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):

View File

@ -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',

View File

@ -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

View File

@ -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 \

View File

@ -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 \