全部產品
Search
文件中心

DataHub:相容Kafka

更新時間:Aug 27, 2026

DataHub已經完全相容Kafka協議,支援使用原生Kafka用戶端對DataHub進行讀寫操作。

Kafka映射DataHub介紹

Topic類型

Kafka 與 DataHub 的 Topic 擴縮容機制不同。建立 DataHub Topic 時,需將擴容模式設定為擴充模式以適配 Kafka 的行為。擴充模式下,Topic 不支援 Shard 的分裂與合併作業,僅支援增加 Shard,暫不支援減少 Shard。

Topic命名

Kafka Topic 映射為 DataHub 的 Project 和 Topic,兩者以英文句號(.)分隔。映射規則如下:

  • 以首個.為界,左側為 DataHub Project,右側為 DataHub Topic。例如 test_project.test_topic 映射為 Project test_project、Topic test_topic。

  • 若名稱中包含多個 .,僅首個 . 作為分隔字元,其餘 . 和 - 均替換為 _。

Partition

DataHub 中每個處於 Active 狀態的 Shard 對應 Kafka 的 1 個 Partition。例如,當前有 5 個 Active Shard,則等效於 Kafka 有 5 個 Partition,寫入時可指定 Partition 範圍為 [0, 4]。未指定時,由 Kafka 用戶端自動分配 Partition。

Tuple Topic

從 Kafka 向 Tuple Topic 寫入資料時,Topic Schema 必須為 1 列或 2 列,且類型均為 STRING,否則寫入失敗。

  • 列數為 1 時僅寫入 Value,Key 資料將被丟棄;

  • 列數為 2 時,第 1 列和第 2 列分別對應 Key 和 Value。

此外,Tuple Topic 不支援寫入位元據,否則會出現亂碼,位元據請寫入 Blob Topic。

Blob Topic

從 Kafka 向 Blob Topic 寫入資料時,Kafka 訊息的 Value 寫入 Blob 欄位。如果訊息的 Key 不為 NULL,則將其寫入 DataHub 的 Attribute,屬性名稱固定為 __kafka_key__,屬性值為 Kafka 訊息的 Key。

Header

Kafka 的 Header 對應 DataHub 的 Attribute。如果 Header 的 Value 為 NULL,則跳過該 Header,不寫入 Attribute。建議不要使用 __kafka_key__ 作為 Header 的 Key,以避免與 Blob Topic 中的內建屬性名稱衝突。

Consumer Group

DataHub 的消費組即訂閱 ID,僅支援訂閱單個 Topic;Kafka 的 Consumer Group 可同時訂閱多個 Topic。為相容 Kafka 的訂閱者式,DataHub 提供了 Group 功能:使用者可在 Project 下建立 Group 並綁定多個 Topic,通過該 Group 統一訂閱同一 Project 下的多個 Topic。

Group 在服務端封裝了多個 DataHub 訂閱。綁定 Topic 後,Topic 頁面的訂閱列表中會出現由 Group 自動建立的訂閱記錄。請勿手動刪除該訂閱,否則 Group 將無法繼續訂閱該 Topic,且已有的消費點位會丟失。

單個 Group 最多支援訂閱 50 個 Topic,如需訂閱更多,請提交工單申請。

Kafka配置參數

C=Consumer, P=Producer, S=Streams

參數

C/P/S

可選配置

是否必須

描述

bootstrap.servers

*

參考Kafka網域名稱列表

是

security.protocol

*

SASL_SSL

是

為了保證資料轉送的安全性,Kafka寫入DataHub預設使用SSL加密傳輸

sasl.mechanism

*

PLAIN

是

AK認證方式,僅支援PLAIN

compression.type

P

LZ4

否

是否開啟壓縮傳輸,目前僅支援LZ4

group.id

C

project.topic:subId

或者

project.group

是

使用project.topic:subId時必須和訂閱的topic保持一致,否則無法讀取資料,推薦使用project.group

partition.assignment.strategy

C

org.apache.kafka.clients.consumer.RangeAssignor

否

Kafka預設為RangeAssignor,並且DataHub目前只支援RangeAssignor,請不要修改此配置

session.timeout.ms

C/S

[60000, 180000]

否

kafka預設為10000, 但是因為DataHub限制最小為60000,所以這裡預設會變為60000

heartbeat.interval.ms

C/S

建議session.timeout.ms的 2/3

否

Kafka預設為3000,但是因為session.timeout.ms會被預設修改為60000,所以這裡建議顯示設定為40000,否則heartbeat請求會過於頻繁

application.id

S

project.topic:subId

或

project.group

是

使用project.topic:subId時必須和訂閱的topic保持一致,否則無法讀取資料,推薦使用project.group

以上列出的是使用 Kafka 用戶端寫入 DataHub 時需要重點關注的參數。其餘用戶端參數(如 retries、batch.size)的行為與原生 Kafka 一致,不受影響。服務端參數不會改變 DataHub 的實際行為,例如無論 acks 設定為何值,DataHub 均在資料完全寫入成功後才返回確認。

Kafka網域名稱列表

地區

Region

外網Endpoint

傳統網路ECS Endpoint

VPC ECS Endpoint

華東1(杭州)

cn-hangzhou

dh-cn-hangzhou.aliyuncs.com:9092

dh-cn-hangzhou.aliyun-inc.com:9093

dh-cn-hangzhou-int-vpc.aliyuncs.com:9094

華東2(上海)

cn-shanghai

dh-cn-shanghai.aliyuncs.com:9092

dh-cn-shanghai.aliyun-inc.com:9093

dh-cn-shanghai-int-vpc.aliyuncs.com:9094

華北2(北京)

cn-beijing

dh-cn-beijing.aliyuncs.com:9092

dh-cn-beijing.aliyun-inc.com:9093

dh-cn-beijing-int-vpc.aliyuncs.com:9094

華北3(張家口)

cn-zhangjiakou

dh-cn-zhangjiakou.aliyuncs.com:9092

dh-cn-zhangjiakou.aliyun-inc.com:9093

dh-cn-zhangjiakou-int-vpc.aliyuncs.com:9094

華南1(深圳)

cn-shenzhen

dh-cn-shenzhen.aliyuncs.com:9092

dh-cn-shenzhen.aliyun-inc.com:9093

dh-cn-shenzhen-int-vpc.aliyuncs.com:9094

亞太地區東南1(新加坡)

ap-southeast-1

dh-ap-southeast-1.aliyuncs.com:9092

dh-ap-southeast-1.aliyun-inc.com:9093

dh-ap-southeast-1-int-vpc.aliyuncs.com:9094

亞太地區東南3(吉隆坡)

ap-southeast-3

dh-ap-southeast-3.aliyuncs.com:9092

dh-ap-southeast-3.aliyun-inc.com:9093

dh-ap-southeast-3-int-vpc.aliyuncs.com:9094

歐洲中部1(法蘭克福)

eu-central-1

dh-eu-central-1.aliyuncs.com:9092

dh-eu-central-1.aliyun-inc.com:9093

dh-eu-central-1-int-vpc.aliyuncs.com:9094

上海金融雲

cn-shanghai-finance-1

dh-cn-shanghai-finance-1.aliyuncs.com:9092

dh-cn-shanghai-finance-1.aliyun-inc.com:9093

dh-cn-shanghai-finance-1-int-vpc.aliyuncs.com:9094

中國香港

cn-hongkong

dh-cn-hongkong.aliyuncs.com:9092

dh-cn-hongkong.aliyun-inc.com:9093

dh-cn-hongkong-int-vpc.aliyuncs.com:9094

建立Topic樣本

  1. 建立頁面

    注意需要開啟Shard擴充模式。

  2. 建立代碼

    目前不支援通過 Kafka API 建立 Topic,需使用 DataHub SDK 建立,並將 ExpandMode 設定為 ONLY_EXTEND。Maven 依賴版本要求 2.19.0 及以上。

    Access Key 和 Secret Key 推薦通過環境變數配置,避免寫入程式碼在工程代碼中。阿里雲帳號的 AccessKey 擁有所有 API 的存取權限,建議使用 RAM 使用者的 AccessKey 進行 API 訪問或日常營運,以降低 AccessKey 泄露後危及帳號下所有資源的風險。

    datahub.endpoint=<yourEndpoint>
    datahub.accessId=<yourAccessKeyId>
    datahub.accessKey=<yourAccessKeySecret>
    <dependency>
      <groupId>com.aliyun.datahub</groupId>
      <artifactId>aliyun-sdk-datahub</artifactId>
      <version>2.19.0-public</version>
    </dependency>
    @Value("${datahub.endpoint}")
    String endpoint ;
    @Value("${datahub.accessId}")
    String accessId;
    @Value("${datahub.accessKey}")
    String accessKey;
    public class CreateTopic {
        public static void main(String[] args) {
            DatahubClient datahubClient = DatahubClientBuilder.newBuilder()
                    .setDatahubConfig(
                            new DatahubConfig(endpoint,
                                    new AliyunAccount(accessId, accessKey)))
                    .build();
    
            int shardCount = 1;
            int lifeCycle = 7;
    
            try {
                datahubClient.createTopic("test_project", "test_topic", shardCount, lifeCycle, RecordType.BLOB, "comment", ExpandMode.ONLY_EXTEND);
            } catch (DatahubClientException e) {
                e.printStackTrace();
            }
        }
    }

建立Group樣本

  1. 建立頁面

    單擊建立Group,將需要訂閱的topic添加到右側列表。建立完成後仍然可以修改綁定的topic列表。在topic的訂閱列表頁面可以看到group自動建立的訂閱。

  2. 建立代碼

    maven依賴版本需為2.21.6-public或更高版本

    <dependency>
      <groupId>com.aliyun.datahub</groupId>
      <artifactId>aliyun-sdk-datahub</artifactId>
      <version>2.21.6-public</version>
    </dependency>
    @Value("${datahub.endpoint}")
    String endpoint ;
    @Value("${datahub.accessId}")
    String accessId;
    @Value("${datahub.accessKey}")
    String accessKey;
    public class CreateGroup {
        public static void main(String[] args) {
            DatahubClient datahubClient = DatahubClientBuilder.newBuilder()
                    .setDatahubConfig(
                            new DatahubConfig(endpoint,
                                    new AliyunAccount(accessId, accessKey)))
                    .build();
    
            List<String> topicList = new ArrayList<>();
            topicList.add("test_project.topic1");
            topicList.add("test_project.topic2");
            topicList.add("test_project.topic3");
    
            try {
                // 建立kafka group
                datahubClient.createKafkaGroup("test_project", "test_topic", "test comment");
    
                // 將需要訂閱的topic綁定到group上
                datahubClient.updateTopicsForKafkaGroup("test_project", "test_topic", topicList, UpdateKafkaGroupMode.ADD);
            } catch (DatahubClientException e) {
                e.printStackTrace();
            }
        }
    }

Producer樣本

產生kafka_client_producer_jaas.conf檔案

建立檔案kafka_client_producer_jaas.conf,儲存到任意路徑,檔案內容如下。

KafkaClient {
  org.apache.kafka.common.security.plain.PlainLoginModule required
  username="accessId"
  password="accessKey";
};

maven依賴

Kafka-client版本至少大於等於0.10.0.0,推薦2.4.0。

<dependency>
    <groupId>org.apache.kafka</groupId>
    <artifactId>kafka-clients</artifactId>
    <version>2.4.0</version>
</dependency>

範例程式碼

public class ProducerExample {
    static {
        System.setProperty("java.security.auth.login.config", "src/main/resources/kafka_client_producer_jaas.conf");
    }

    public static void main(String[] args) {

        Properties properties = new Properties();
        properties.put("bootstrap.servers", "dh-cn-hangzhou.aliyuncs.com:9092");
        properties.put("security.protocol", "SASL_SSL");
        properties.put("sasl.mechanism", "PLAIN");
        properties.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
        properties.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
        properties.put("compression.type", "lz4");

        String KafkaTopicName = "test_project.test_topic";
        Producer<String, String> producer = new KafkaProducer<String, String>(properties);

        try {
            List<Header> headers = new ArrayList<>();
            RecordHeader header1 = new RecordHeader("key1", "value1".getBytes());
            RecordHeader header2 = new RecordHeader("key2", "value2".getBytes());
            headers.add(header1);
            headers.add(header2);

            ProducerRecord<String, String> record = new ProducerRecord<>(KafkaTopicName, 0, "key", "Hello DataHub!", headers);

            // sync send
            producer.send(record).get();

        } catch (InterruptedException e) {
            e.printStackTrace();
        } catch (ExecutionException e) {
            e.printStackTrace();
        } finally {
            producer.close();
        }
    }
}

運行結果

運行成功之後,抽樣確認結果是否正確。

Consumer樣本

產生kafka_client_producer_jaas.conf檔案和maven依賴參考Producer樣本。

新加入的consumer需要十幾秒左右分配shard,分配完成後即可消費。

範例程式碼

使用kafka group樣本(推薦)

package com.aliyun.datahub.kafka.demo;

import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaConsumer;

import java.time.Duration;
import java.util.ArrayList;
import java.util.List;
import java.util.Properties;
public class ConsumerExample2 {

    static {
        System.setProperty("java.security.auth.login.config", "src/main/resources/kafka_client_producer_jaas.conf");
    }

    public static void main(String[] args) {

        Properties properties = new Properties();
        properties.put("bootstrap.servers", "dh-cn-hangzhou.aliyuncs.com:9092");
        properties.put("security.protocol", "SASL_SSL");
        properties.put("sasl.mechanism", "PLAIN");
        // group.id填project.group
        properties.put("group.id", "test_project.test_kafka_group");
        properties.put("auto.offset.reset", "earliest");
        properties.put("session.timeout.ms", "60000");
        properties.put("heartbeat.interval.ms", "40000");
        properties.put("ssl.endpoint.identification.algorithm", "");
        properties.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
        properties.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
        KafkaConsumer<String, String> kafkaConsumer = new KafkaConsumer<String, String>(properties);

        List<String> topicList = new ArrayList<>();
        topicList.add("test_project.test_topic1");
        topicList.add("test_project.test_topic2");
        topicList.add("test_project.test_topic3");
        // 使用kafka group可以同時訂閱多個topic
        kafkaConsumer.subscribe(topicList);

        while (true) {
            ConsumerRecords<String, String> records = kafkaConsumer.poll(Duration.ofSeconds(5));

            for (ConsumerRecord<String, String> record : records) {
                System.out.println(record.toString());
            }
        }
    }
}

使用project.topic.subid樣本

package com.aliyun.datahub.kafka.demo;

import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaConsumer;

import java.time.Duration;
import java.util.Collections;
import java.util.Properties;
public class ConsumerExample {

    static {
        System.setProperty("java.security.auth.login.config", "src/main/resources/kafka_client_producer_jaas.conf");
    }

    public static void main(String[] args) {

        Properties properties = new Properties();
        properties.put("bootstrap.servers", "dh-cn-hangzhou.aliyuncs.com:9092");
        properties.put("security.protocol", "SASL_SSL");
        properties.put("sasl.mechanism", "PLAIN");
        // group.id填project.topic.subId
        properties.put("group.id", "test_project.test_topic:1611039998153N71KM");
        properties.put("auto.offset.reset", "earliest");
        properties.put("session.timeout.ms", "60000");
        properties.put("heartbeat.interval.ms", "40000");
        properties.put("ssl.endpoint.identification.algorithm", "");
        properties.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
        properties.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
        KafkaConsumer<String, String> kafkaConsumer = new KafkaConsumer<String, String>(properties);

        // 使用project.topic.subId的方式只能訂閱單個topic
        kafkaConsumer.subscribe(Collections.singletonList("test_project.test_topic"));

        while (true) {
            ConsumerRecords<String, String> records = kafkaConsumer.poll(Duration.ofSeconds(5));

            for (ConsumerRecord<String, String> record : records) {
                System.out.println(record.toString());
            }
        }
    }
}

運行結果

運行成功之後,便可以在終端看到讀取到的資料。

ConsumerRecord(topic = test_project.test_topic, partition = 0, leaderEpoch = 0, offset = 0, LogAppendTime = 1611040892661, serialized key size = 3, serialized value size = 14, headers = RecordHeaders(headers = [RecordHeader(key = key1, value = [118, 97, 108, 117, 101, 49]), RecordHeader(key = key2, value = [118, 97, 108, 117, 101, 50])], isReadOnly = false), key = key, value = Hello DataHub!)

本樣本中同一個請求返回的資料的LogAppendTime是相同的,是該請求返回所有的資料的寫入DataHub時間的最大值。

Streams樣本

maven依賴

<dependency>
    <groupId>org.apache.kafka</groupId>
    <artifactId>kafka-clients</artifactId>
    <version>2.4.0</version>
</dependency>
<dependency>
    <groupId>org.apache.kafka</groupId>
    <artifactId>kafka-streams</artifactId>
    <version>2.4.0</version>
</dependency>

程式碼範例

這裡讀取test_project下input的資料,將key和value的字串轉為小寫重新寫入output。

public class StreamExample {

    static {
        System.setProperty("java.security.auth.login.config", "src/main/resources/kafka_client_producer_jaas.conf");
    }

    public static void main(final String[] args) {
        final String input = "test_project.input";
        final String output = "test_project.output";
        final Properties properties = new Properties();
        properties.put("bootstrap.servers", "dh-cn-hangzhou.aliyuncs.com:9092");
        properties.put("application.id", "test_project.input:1611293595417QH0WL");
        properties.put("security.protocol", "SASL_SSL");
        properties.put("sasl.mechanism", "PLAIN");
        properties.put("session.timeout.ms", "60000");
        properties.put("heartbeat.interval.ms", "40000");
        properties.put("auto.offset.reset", "earliest");

        final StreamsBuilder builder = new StreamsBuilder();
        TestMapper testMapper = new TestMapper();
        builder.stream(input, Consumed.with(Serdes.String(), Serdes.String()))
                .map(testMapper)
                .to(output, Produced.with(Serdes.String(), Serdes.String()));

        final KafkaStreams streams = new KafkaStreams(builder.build(), properties);
        final CountDownLatch latch = new CountDownLatch(1);

        Runtime.getRuntime().addShutdownHook(new Thread("streams-shutdown-hook") {
            @Override
            public void run() {
                streams.close();
                latch.countDown();
            }
        });

        try {
            streams.start();
            latch.await();
        } catch (final Throwable e) {
            System.exit(1);
        }
        System.exit(0);
    }

    static class TestMapper implements KeyValueMapper<String, String, KeyValue<String, String>> {

        @Override
        public KeyValue<String, String> apply(String s, String s2) {
            return new KeyValue<>(StringUtils.lowerCase(s), StringUtils.lowerCase(s2));
        }
    }
}

運行結果

啟動Streams任務之後,分配shard大概需要1分鐘左右,1分鐘之後就可以在控制台看到當前的task數量,task數量和輸入topic的shard數量保持一致,樣本輸入topic為3個shard。

currently assigned active tasks: [0_0, 0_1, 0_2]
currently assigned standby tasks: []
revoked active tasks: []
  revoked standby tasks: []

shard分配成功之後,可以向input中寫入一組測試資料 (AAAA,BBBB),(CCCC,DDDD),(EEEE,FFFF),再output抽樣查看資料是否正確寫入。

注意事項

  • 目前不支援事務、等冪。

  • 目前Kafka用戶端無法自動建立DataHub Topic,寫入之前需要保證已建立Topic。

  • Consumer目前最多隻能訂閱一個topic。

  • Consumer讀取的資料時間戳記均為LogAppendTime,表示DataHub的落盤時間,單個請求返回的所有資料時間戳記相同,為所有資料時間戳記的最大值,所以如果讀取的時間戳記可能會大於實際的落盤時間

  • Streams輸入topic目前僅支援一個,輸出可以多個topic。

  • Streams目前只支援無狀態的任務。

  • 支援Kafka版本為0.10.0 -> 2.4.0。

常見問題

寫入資料時串連斷開

Selector - [Producer clientId=producer-1] Connection with dh-cn-shenzhen.aliyuncs.com disconnected
java.io.EOFException
    at org.apache.kafka.common.network.SslTransportLayer.read(SslTransportLayer.java:573)
    ...

Kafka 的中繼資料(Metadata)請求與資料寫入請求使用不同的串連。

用戶端首先建立串連擷取中繼資料,再根據返回的 Broker 資訊建立第二個串連用於寫入資料,後續所有請求均通過第二個串連發送。

第一個串連因此處於閑置狀態,服務端會在串連閑置逾時後主動關閉,此時可能產生串連斷開的錯誤記錄檔。如果資料寫入正常,可忽略該錯誤。

啟動kafka用戶端失敗

Caused by: org.apache.kafka.common.errors.SslAuthenticationException: SSL handshake failed
Caused by: javax.net.ssl.SSLHandshakeException: No subject alternative names matching IP address 100.67.134.161 found

添加配置properties.put("ssl.endpoint.identification.algorithm", "");。

Consumer消費過程中出現DisconnectException

[INFO][Consumer clientId=client-id, groupId=consumer-project.topic:subid] Error sending fetch request (sessionId=INVALID, epoch=INITIAL) to node 1: {}.
org.apache.kafka.common.errors.DisconnectException

Kafka的用戶端需要與服務端保持TCP長串連,一般情況是因為網路抖動造成的,用戶端有重試邏輯,因此不會對用戶端的消費造成影響。