All Products
Search
Document Center

DataHub:Kafka compatibility

Last Updated:Aug 26, 2026

DataHub is fully compatible with the Apache Kafka protocol. You can use native Kafka clients to read data from and write data to DataHub.

Kafka to DataHub mapping

Topic types

Kafka and DataHub have different mechanisms for scaling topics. To ensure compatibility with Kafka's behavior, you must set the scaling mode to ONLY_EXTEND when you create a DataHub topic. In this mode, you can only add new shards to a topic. This mode does not allow splitting or merging shards, and shard removal is not yet supported.

Topic naming

A Kafka topic name maps to a DataHub project and topic, separated by a period (.). The mapping follows these rules:

  • The part before the first. is the DataHub project, and the part after it is the DataHub topic. For example, test_project.test_topic maps to thetest_project project and thetest_topic topic.

  • If a name contains multiple. characters, only the first. acts as a separator. All remaining. and- characters are replaced with_.

Partitions

Each active shard in DataHub corresponds to one partition in Kafka. For example, if a topic has five active shards, it is equivalent to a Kafka topic with five partitions. When you write data, you can specify a partition ID in the range [0, 4]. If you do not specify a partition, the Kafka client automatically assigns one.

Tuple topic

When you write data from Kafka to a Tuple topic, the topic schema must have one or two columns, both of type STRING. Otherwise, the write operation fails.

  • If the schema has one column, only the value is written, and the key is discarded.

  • If the schema has two columns, the first and second columns correspond to the key and value, respectively.

In addition, do not write binary data to a Tuple topic, as this causes garbled characters. To store binary data, use a Blob topic.

Blob topic

When you write data from Kafka to a Blob topic, the Kafka message value is written to the Blob field. If the message key is not NULL, it is written as a DataHub attribute. The attribute name is __kafka_key__, and its value is the Kafka message key.

Headers

Kafka headers map to DataHub attributes. If a header's value is NULL, it is ignored and not written as an attribute. We recommend that you do not use __kafka_key__ as a header key to avoid conflicts with the built-in attribute name in Blob topics.

Consumer groups

In DataHub, a subscription ID acts as a consumer group but can only subscribe to a single topic. In contrast, a Kafka consumer group can subscribe to multiple topics simultaneously. To provide compatibility with Kafka's subscription model, DataHub offers a group feature. You can create a group within a project and bind it to multiple topics, allowing you to subscribe to all of them under a single group.

A group internally manages multiple DataHub subscriptions on the server. After you bind a topic, the group automatically creates a subscription, which appears in the subscription list on the topic's details page. Do not delete this subscription manually. Doing so will prevent the group from subscribing to the topic and will cause all existing consumption offsets to be lost.

A single group can subscribe to a maximum of 50 topics. To subscribe to more topics, submit a ticket.

Kafka parameters

C=Consumer, P=Producer, S=Streams

Parameter

C/P/S

Value

Required

Description

bootstrap.servers

*

See the Kafka endpoints section.

Yes

security.protocol

*

SASL_SSL

Yes

To ensure data security, connections from Kafka to DataHub use SSL encryption by default.

sasl.mechanism

*

PLAIN

Yes

The authentication mechanism for AccessKey credentials. Only PLAIN is supported.

compression.type

P

LZ4

No

Specifies the compression type for messages. Currently, only LZ4 is supported.

group.id

C

project.topic:subId

or

project.group

Yes

If you use the project.topic:subId format, the ID must match the subscribed topic. Otherwise, data cannot be read. We recommend using the project.group format.

partition.assignment.strategy

C

org.apache.kafka.clients.consumer.RangeAssignor

No

The default partition assignment strategy in Kafka is RangeAssignor. DataHub currently supports only this strategy. Do not modify this parameter.

session.timeout.ms

C/S

[60000, 180000]

No

The default in Kafka is 10,000 ms. However, because DataHub requires a minimum of 60,000 ms, this value is automatically adjusted to 60,000 ms.

heartbeat.interval.ms

C/S

Recommended: 2/3 of session.timeout.ms

No

The Kafka default is 3,000 ms. Becausesession.timeout.ms is adjusted to 60,000 ms, werecommend that you explicitly set this value to 40000 to prevent frequent heartbeat requests.

application.id

S

project.topic:subId

or

project.group

Yes

If you use the project.topic:subId format, the ID must match the subscribed topic. Otherwise, data cannot be read. We recommend using the project.group format.

This table lists the key parameters to review when using a Kafka client with DataHub. Other client-side parameters, such as retries and batch.size, behave as they do in native Kafka. Server-side parameters do not change the actual behavior of DataHub. For example, regardless of the value of acks, DataHub returns a confirmation only after data is fully written.

Kafka endpoints

Region

Region ID

Public endpoint

ECS endpoint (classic network)

ECS endpoint (VPC)

China (Hangzhou)

cn-hangzhou

dh-cn-hangzhou.aliyuncs.com:9092

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

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

China (Shanghai)

cn-shanghai

dh-cn-shanghai.aliyuncs.com:9092

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

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

China (Beijing)

cn-beijing

dh-cn-beijing.aliyuncs.com:9092

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

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

China (Zhangjiakou)

cn-zhangjiakou

dh-cn-zhangjiakou.aliyuncs.com:9092

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

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

China (Shenzhen)

cn-shenzhen

dh-cn-shenzhen.aliyuncs.com:9092

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

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

Singapore

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

Malaysia (Kuala Lumpur)

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

Germany (Frankfurt)

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

China East 2 Finance

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

China (Hong Kong)

cn-hongkong

dh-cn-hongkong.aliyuncs.com:9092

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

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

Create a topic

  1. Create a topic in the console

    When you create a topic, enable Shard Expand Mode.

  2. Create a topic by using the SDK

    You cannot create topics by using the Kafka API. You must use the DataHub SDK and set ExpandMode to ONLY_EXTEND. The required Maven dependency version is 2.19.0 or later.

    We recommend that you use environment variables to configure your AccessKey ID and AccessKey Secret and avoid hard-coding them into your project code. An Alibaba Cloud account's AccessKey pair has permissions for all API operations. For better security, use an AccessKey pair from a RAM user for API access or daily operations to reduce the risk of credential leaks.

    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();
            }
        }
    }

Create a group

  1. Create a group in the console

    Click Create Group, and then add the topics you want to subscribe to from the list on the right. You can modify the bound topics after the group is created. The group automatically creates a subscription, which appears on the topic's subscription list page.

  2. Create a group by using the SDK

    The Maven dependency version must be 2.21.6-public or later.

    <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 {
                // Create a Kafka group.
                datahubClient.createKafkaGroup("test_project", "test_topic", "test comment");
    
                // Bind the topics to the group for subscription.
                datahubClient.updateTopicsForKafkaGroup("test_project", "test_topic", topicList, UpdateKafkaGroupMode.ADD);
            } catch (DatahubClientException e) {
                e.printStackTrace();
            }
        }
    }

Producer example

The kafka_client_producer_jaas.conf file

Create a file named kafka_client_producer_jaas.conf in any directory and add the following content.

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

Maven dependency

The Kafka client version must be 0.10.0.0 or later. We recommend version 2.4.0.

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

Sample code

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();
        }
    }
}

Result

After the code runs successfully, you can sample data to verify the result.

Consumer example

For information about how to generate thekafka_client_producer_jaas.conf file and add the Maven dependency, see the producer example.

When a new consumer joins, shard assignment takes 10 to 20 seconds. After the assignment is complete, the consumer can start consuming data.

Sample code

Using a Kafka group (recommended)

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");
        // Set group.id to the project.group format.
        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");
        // By using a Kafka group, you can subscribe to multiple topics.
        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());
            }
        }
    }
}

Using 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");
        // Set group.id to the project.topic:subId format.
        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);

        // When using the project.topic:subId format, you can only subscribe to a single 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());
            }
        }
    }
}

Result

After the code runs successfully, you can see the consumed data in your terminal.

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!)

In this example, all data records returned in a single request have the same LogAppendTime, which is the latest timestamp among all records in that batch.

Streams example

Maven dependency

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

Sample code

This example reads data from an input topic within test_project, converts the key and value strings to lowercase, and writes the result to an output topic.

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));
        }
    }
}

Result

After you start the Streams task, shard assignment takes about one minute. After that, you can see the number of current tasks in the console. The number of tasks matches the number of shards in the input topic. In this example, the input topic has three shards.

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

After the shards are assigned, you can write test data such as (AAAA,BBBB),(CCCC,DDDD),(EEEE,FFFF) to the input topic. Then, sample data from the output topic to verify that it was written correctly.

Usage notes

  • Transactions and idempotency are not supported.

  • Kafka clients cannot automatically create topics in DataHub. You must create the topic before writing data to it.

  • When using a subscription ID (project.topic:subid) as the group.id, a consumer can subscribe to only one topic. To subscribe to multiple topics, use a DataHub group.

  • The timestamp for data read by a consumer is always the LogAppendTime, which indicates when the data was written to DataHub. All records in a single fetch request share the same timestamp: the latest timestamp in that batch. This means the read timestamp might be later than the actual write time.

  • A Streams application supports only one input topic but can have multiple output topics.

  • Only stateless Streams tasks are supported.

  • Supported Kafka versions range from 0.10.0 to 2.4.0.

FAQ

Connection disconnects when writing data

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 requests and data write requests use different connections.

The client first establishes a connection to fetch metadata. It then uses the returned broker information to establish a second connection for writing data. All subsequent requests are sent over this second connection.

The first connection, now idle, is automatically closed by the server after a timeout. This may generate a disconnection error in the logs. You can ignore this error if data is being written successfully.

Kafka client fails to start

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

Add the following property to your configuration: properties.put("ssl.endpoint.identification.algorithm", "");.

DisconnectException during consumption

[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

The Kafka client must maintain a persistent TCP connection with the server. This exception is usually caused by network jitter. The client has built-in retry logic, so this error typically does not affect consumption.