脚本能自动删除旧Kafka主题吗?

wen 实用脚本 4

本文目录导读:

脚本能自动删除旧Kafka主题吗?

  1. 配置准备
  2. 脚本实现
  3. 定时任务配置
  4. 注意事项
  5. 推荐做法

是的,可以通过脚本自动删除旧的Kafka主题,但需要注意,Kafka默认禁用了自动删除主题的功能,需要先进行配置。

配置准备

启用主题删除功能

在Kafka服务器配置文件 server.properties 中:

# 设置为true以启用删除功能
delete.topic.enable=true

配置文件路径

根据Kafka版本不同,配置文件位置可能不同:

  • Confluent版本:/etc/kafka/server.properties
  • Apache Kafka源码版:$KAFKA_HOME/config/server.properties

修改后需重启Kafka服务。

脚本实现

Bash脚本示例(基于时间判断)

#!/bin/bash
# Kafka配置
KAFKA_HOME="/path/to/kafka"
BOOTSTRAP_SERVER="localhost:9092"
# 删除3天前的主题
RETENTION_DAYS=3
# 匹配旧主题的关键字(可选)
TOPIC_PATTERN="^old-"
# 计算时间戳
CUTOFF_DATE=$(date -d "$RETENTION_DAYS days ago" +%s)
# 获取所有主题列表
topics=$(${KAFKA_HOME}/bin/kafka-topics.sh --bootstrap-server ${BOOTSTRAP_SERVER} --list)
echo "开始扫描旧主题..."
for topic in $topics; do
    # 跳过系统主题(__开头)
    if [[ "$topic" == __* ]]; then
        continue
    fi
    # 可选:匹配特定模式
    if [[ -n "$TOPIC_PATTERN" ]] && [[ ! "$topic" =~ $TOPIC_PATTERN ]]; then
        continue
    fi
    # 获取主题的创建时间(需要Kafka 2.8+支持)
    # 方法1:通过describe获取
    create_time=$(${KAFKA_HOME}/bin/kafka-topics.sh --bootstrap-server ${BOOTSTRAP_SERVER} \
        --describe --topic "$topic" 2>/dev/null | grep "Topic: $topic" | awk '{print $NF}')
    # 方法2:通过Zookeeper获取(旧版本Kafka)
    # create_time=$(echo stat /brokers/topics/${topic} | zookeeper-server localhost:2181 | grep "created" | awk '{print $2}')
    # 检查时间戳
    if [[ -n "$create_time" ]] && [[ "$create_time" -lt "$CUTOFF_DATE"  ]]; then
        echo "删除主题: $topic (创建时间: $(date -d @$create_time '+%Y-%m-%d %H:%M:%S'))"
        ${KAFKA_HOME}/bin/kafka-topics.sh --bootstrap-server ${BOOTSTRAP_SERVER} \
            --delete --topic "$topic" 2>/dev/null
        if [ $? -eq 0 ]; then
            echo "  删除成功"
        else
            echo "  删除失败,请检查权限或配置"
        fi
    fi
done

替代方案:基于主题元数据

如果不想依赖创建时间,可以基于主题的其他元数据规则:

#!/bin/bash
# 删除没有消费者组的主题(无人使用)
KAFKA_HOME="/path/to/kafka"
BOOTSTRAP_SERVER="localhost:9092"
DRY_RUN=true  # 设为false执行实际删除
all_topics=$(${KAFKA_HOME}/bin/kafka-topics.sh --bootstrap-server ${BOOTSTRAP_SERVER} --list)
# 获取有消费者组的主题
active_topics=$(${KAFKA_HOME}/bin/kafka-consumer-groups.sh \
    --bootstrap-server ${BOOTSTRAP_SERVER} --all-groups --list 2>/dev/null)
for topic in $all_topics; do
    # 跳过系统主题
    if [[ "$topic" == __* ]]; then
        continue
    fi
    # 检查主题是否活跃(有消费者消费)
    if ! echo "$active_topics" | grep -q "$topic"; then
        if [ "$DRY_RUN" = true ]; then
            echo "[DRY RUN] 将删除主题: $topic"
        else
            echo "删除主题: $topic"
            ${KAFKA_HOME}/bin/kafka-topics.sh --bootstrap-server ${BOOTSTRAP_SERVER} \
                --delete --topic "$topic"
        fi
    fi
done

定时任务配置

使用crontab实现定期清理:

# 每天凌晨2点执行
0 2 * * * /path/to/delete_old_topics.sh >> /var/log/kafka_topic_cleanup.log 2>&1

或者使用systemd定时器(Linux系统):

# /etc/systemd/system/kafka-topic-cleanup.timer
[Unit]
Description=Kafka旧主题清理定时器
[Timer]
OnCalendar=daily
Persistent=true
[Install]
WantedBy=timers.target

注意事项

  1. 备份配置:在自动化删除前,先手动测试一次
  2. 日志记录:记录所有删除操作以备审计
  3. 数据安全:删除操作不可逆,确保业务不受影响
  4. 高版本Kafka:Kafka 2.8+支持获取主题创建时间;旧版本需要通过Zookeeper或自定义元数据判断
  5. 幂等性:添加--if-exists参数避免错误
  6. 权限验证:如果启用了SASL/SSL认证,脚本中需要添加认证参数

推荐做法

更稳妥的方式是使用主题生命周期管理工具,如:

这样可以在删除前发送告警、设置白名单等,避免误删。

抱歉,评论功能暂时关闭!