有没有办法在每次运行之前删除主题中的所有数据或删除主题?

Tom*_*myT 77 apache-kafka apache-zookeeper

有没有办法在每次运行之前删除主题中的所有数据或删除主题?

我可以修改KafkaConfig.scala文件来更改logRetentionHours属性吗?一旦消费者阅读消息,是否有消息被删除的方式?

我正在使用生产者从某个地方获取数据并将数据发送到消费者消费的特定主题,我可以在每次运行时删除该主题中的所有数据吗?我只想在主题中每次都有新数据.有没有办法以某种方式重新初始化该主题?

Pat*_*ick 61

正如我在这里提到的Purge Kafka Queue:

在Kafka 0.8.2中测试,用于快速入门示例:首先,在config文件夹下的server.properties文件中添加一行:

delete.topic.enable=true
Run Code Online (Sandbox Code Playgroud)

然后,您可以运行此命令:

bin/kafka-topics.sh --zookeeper localhost:2181 --delete --topic test
Run Code Online (Sandbox Code Playgroud)

  • 这应该是恕我直言的答案。 (2认同)
  • 顺便说一句,添加选项后你不需要重启Kafka服务器,万一有人想知道. (2认同)

Hil*_*ild 53

不要认为它还支持.看看这个JIRA问题 "添加删除主题支持".

要手动删除:

  1. 关闭群集
  2. 清理kafka日志目录(由log.dirkafka 配置文件中的属性指定)以及zookeeper数据
  3. 重新启动群集

对于任何给定的主题,您可以做的是

  1. 停止卡夫卡
  2. 具体到分区清洁卡夫卡日志,卡夫卡存储其日志文件中的"LOGDIR /主题分区"的格式,所以名为"MyTopic"主题中的日志分区ID 0将被存储在/tmp/kafka-logs/MyTopic-0其中/tmp/kafka-logs被指定的log.dir属性
  3. 重启kafka

这是NOT一个很好的推荐方法,但它应该有效.在Kafka代理配置文件中,该log.retention.hours.per.topic属性用于定义The number of hours to keep a log file before deleting it for some specific topic

此外,消费者在阅读消息后是否有消息被删除的方式?

来自Kafka文档:

Kafka群集保留所有已发布的消息 - 无论它们是否已被消耗 - 在可配置的时间段内.例如,如果将日志保留设置为两天,那么在发布消息后的两天内,它可供消费,之后将被丢弃以释放空间.Kafka的性能在数据大小方面实际上是恒定的,因此保留大量数据不是问题.

实际上,基于每个消费者保留的唯一元数据是消费者在日志中的位置,称为"偏移".这种偏移由消费者控制:消费者通常在读取消息时线性地提升其偏移量,但实际上该位置由消费者控制并且它可以按照其喜欢的任何顺序消费消息.例如,消费者可以重置为较旧的偏移量以进行重新处理.

他们说,为了找到Kafka 0.8 Simple Consumer示例中要读取的起始偏移量

Kafka包含两个常量来帮助,kafka.api.OffsetRequest.EarliestTime()在日志中找到数据的开头并从那里开始流式传输,kafka.api.OffsetRequest.LatestTime()只会传输新的消息.

您还可以在那里找到用于管理消费者端偏移量的示例代码.

    public static long getLastOffset(SimpleConsumer consumer, String topic, int partition,
                                 long whichTime, String clientName) {
    TopicAndPartition topicAndPartition = new TopicAndPartition(topic, partition);
    Map<TopicAndPartition, PartitionOffsetRequestInfo> requestInfo = new HashMap<TopicAndPartition, PartitionOffsetRequestInfo>();
    requestInfo.put(topicAndPartition, new PartitionOffsetRequestInfo(whichTime, 1));
    kafka.javaapi.OffsetRequest request = new kafka.javaapi.OffsetRequest(requestInfo, kafka.api.OffsetRequest.CurrentVersion(),clientName);
    OffsetResponse response = consumer.getOffsetsBefore(request);

    if (response.hasError()) {
        System.out.println("Error fetching data Offset Data the Broker. Reason: " + response.errorCode(topic, partition) );
        return 0;
    }
    long[] offsets = response.offsets(topic, partition);
    return offsets[0];
}
Run Code Online (Sandbox Code Playgroud)

  • 该主题仍将显示在此处,因为它已在zookeeper中列出.你必须递归删除`broker/topics/<topic_to_delete>`下的所有内容以及日志来删除它. (4认同)
  • 更新:从kafka 0.8.2开始,命令更改为:`kafka-run-class.sh kafka.admin.TopicCommand --delete --topic [topic_to_delete] --zookeeper localhost:2181` (4认同)
  • 根据问题链接,您可以删除版本0.8.1之后的主题.您可以通过`kafka-run-class.sh kafka.admin.DeleteTopicCommand`查看详细帮助. (3认同)

Swa*_*shi 12

用kafka 0.10测试

1. stop zookeeper & Kafka server,
2. then go to 'kafka-logs' folder , there you will see list of kafka topic folders, delete folder with topic name
3. go to 'zookeeper-data' folder , delete data inside that.
4. start zookeeper & kafka server again.
Run Code Online (Sandbox Code Playgroud)

注意:如果要删除kafka-logs中的主题文件夹,而不是从zookeeper-data文件夹中删除主题文件夹,那么您将看到主题仍然存在.


小智 6

以下是用于清空和删除Kafka主题的脚本,假设localhost为zookeeper服务器,并且Kafka_Home设置为安装目录:

下面的脚本将通过将其保留时间设置为1秒然后删除配置来清空主题:

#!/bin/bash
echo "Enter name of topic to empty:"
read topicName
/$Kafka_Home/bin/kafka-configs --zookeeper localhost:2181 --alter --entity-type topics --entity-name $topicName --add-config retention.ms=1000
sleep 5
/$Kafka_Home/bin/kafka-configs --zookeeper localhost:2181 --alter --entity-type topics --entity-name $topicName --delete-config retention.ms
Run Code Online (Sandbox Code Playgroud)

要完全删除主题,您必须停止任何适用的kafka代理并从kafka日志目录中删除它的目录(默认值:/ tmp/kafka-logs),然后运行此脚本以从zookeeper中删除该主题.要验证它已从zookeeper中删除,ls/brokers/topics的输出不应再包含主题:

#!/bin/bash
echo "Enter name of topic to delete from zookeeper:"
read topicName
/$Kafka_Home/bin/zookeeper-shell localhost:2181 <<EOF
rmr /brokers/topics/$topicName
ls /brokers/topics
quit
EOF
Run Code Online (Sandbox Code Playgroud)

  • 我想编辑答案,因为第一个命令中有一个小错误.但是不允许进行一个字符编辑.实际上它不是`--add config`而是`--add-config` (2认同)

Iva*_*hov 5

作为一种肮脏的解决方法,您可以调整每个主题的运行时保留设置,例如bin/kafka-topics.sh --zookeeper localhost:2181 --alter --topic my_topic --config retention.bytes=1(retention.bytes = 0也可能有效)

过了一会儿,卡夫卡应该释放空间.与重新创建主题相比,不确定这是否有任何影响.

PS.一旦kafka完成清洁,最好将保留设置恢复.

您还可以使用retention.ms来保存历史数据


Dan*_*n M 5

我们尝试了其他答案所描述的中等程度的成功.真正适合我们的是(Apache Kafka 0.8.1)是类命令

sh kafka-run-class.sh kafka.admin.DeleteTopicCommand --topic yourtopic --zookeeper localhost:2181

  • 尝试0.8.2.1(自制软件),它给出了这个错误.`错误:无法找到或加载主类kafka.admin.DeleteTopicCommand` (7认同)
  • 在0.8.1中尝试过这个.该命令返回"删除成功!" 但是它不会删除日志文件夹中的分区. (2认同)
  • 从新的kafka(0.8.2)开始,它是sh kafka-run-class.sh kafka.admin.TopicCommand --delete --topic [topic_for_delete] --zookeeper localhost:2181。确保delete.topic.enable为true。 (2认同)

Mat*_*ipe 5

对于 brew 用户

如果您brew像我一样使用并浪费了大量时间搜索臭名昭著的kafka-logs文件夹,请不要再担心。(请让我知道这是否适合您和多个不同版本的 Homebrew、Kafka 等 :))

您可能会在以下位置找到它:

地点:

/usr/local/var/lib/kafka-logs


如何真正找到那条路径

(这对您通过 brew 安装的基本上每个应用程序也很有帮助)

1) brew services list

kafka 启动 matbhz /Users/matbhz/Library/LaunchAgents/homebrew.mxcl.kafka.plist

2)打开并阅读plist您在上面找到的内容

3)找到定义server.properties位置的行打开它,在我的例子中:

  • /usr/local/etc/kafka/server.properties

4)寻找log.dirs线路:

log.dirs=/usr/local/var/lib/kafka-logs

5) 前往该位置并删除您想要的主题的日志

6)重启Kafka brew services restart kafka