【问题标题】:Kafka doesn't compress in snappy卡夫卡不会快速压缩
【发布时间】:2015-10-20 16:34:29
【问题描述】:

我很好奇,我想检查一下 snappy 压缩是否适用于 java Kafka 客户端。

为了处理这个问题,我设置了一个小程序。该程序生成 1024 条消息可读数据。它们的大小为 1024 字节。我在树新主题上发送这些消息,然后直接在代理文件系统上检查这些主题的大小。

你可以通过以下代码找到这个程序:

package unit_test.testCompress;

import java.util.HashMap;
import java.util.Map;
import java.util.Random;
import java.util.concurrent.Future;

import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.clients.producer.RecordMetadata;


/**
 * Can be use in order to execute some unit test on compression  
 */
public class TestCompress {

    public static void compress(String type, String version){
        Map<String,Object> configs = new HashMap<String,Object>();
        configs.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
        configs.put("producer.type", "async");
        configs.put("compression.type", type);
        configs.put("value.serializer", "org.apache.kafka.common.serialization.ByteArraySerializer");
        configs.put("partitioner.class", "com.kafkaproducer.RecordPartitioner");
        configs.put("bootstrap.servers", "kafka:9092");


        KafkaProducer<String, byte[]> producer = new KafkaProducer<String, byte[]>(configs);

        Random r = new Random(15415485);
        int size = 1024; //1 Ko
        byte[] buffer = new byte[size];
        for(int i = 0; i < size; i++){
            buffer[i] = (byte) ('A' + (r.nextInt() % 26));
        }
        buffer[size-1] = 0;
        //System.out.println(new String(buffer));
        for(int i = 0; i < size; i++ ){
            Future<RecordMetadata> result = producer.send( new ProducerRecord<String, byte[]>("unit_test_compress_"+version+ "_" + type , buffer));
        }

        producer.close();
    }

    public static void main(String[] args) {

        String version = "v10";
        compress("snappy",version);
        compress("gzip",version);
        compress("none",version);

    }

}

我正在使用以下 maven pom 文件编译此代码:

    <project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
  xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
  <modelVersion>4.0.0</modelVersion>

  <groupId>unit_test</groupId>
  <artifactId>testCompress</artifactId>
  <version>0.0.1-SNAPSHOT</version>
  <packaging>jar</packaging>

  <name>testCompress</name>
  <url>http://maven.apache.org</url>

  <properties>
    <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
  </properties>

  <dependencies>  
     <dependency>
        <groupId>org.apache.kafka</groupId>
        <artifactId>kafka_2.10</artifactId>
        <version>0.8.2.2</version>      
    </dependency>
  </dependencies>
</project>

这个程序在我的电脑上运行得很好。

但是当我直接在我的 kafka 代理上检查结果时,它会给出以下输出:

我认为这意味着 snappy 主题没有压缩(但 gzip 压缩效果很好)。我检查了 vi 文件

我在 Kafka 8.2.1 上知道这个问题:https://issues.apache.org/jira/browse/KAFKA-2189 但我在生产者上使用 Kafka 8.2.2,在经纪人上使用 8.2.1。我也检查了 Snappy 的依赖关系。我正在使用 1.1.1.7

你知道如何在 Kafak 上启用 snappy 压缩吗?我是否忘记了在 kafka 上启用快速压缩的参数?

【问题讨论】:

    标签: java maven apache-kafka


    【解决方案1】:

    在 Kafka ML 上交换后,问题是我的 Kafka Broker 必须升级到 8.2.2 版本。它解决了我的问题。

    【讨论】:

      猜你喜欢
      • 2023-03-22
      • 1970-01-01
      • 2019-01-22
      • 2015-12-14
      • 2017-10-02
      • 1970-01-01
      • 2018-09-15
      • 2014-10-19
      • 1970-01-01
      相关资源
      最近更新 更多