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,但是因为 |
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示例
创建页面
注意需要开启Shard扩展模式。
创建代码
目前不支持通过 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示例
创建页面
单击新建Group,将需要订阅的topic添加到右侧列表。创建完成后仍然可以修改绑定的topic列表。在topic的订阅列表页面可以看到group自动创建的订阅。
创建代码
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.DisconnectExceptionKafka的客户端需要与服务端保持TCP长连接,一般情况是因为网络抖动造成的,客户端有重试逻辑,因此不会对客户端的消费造成影响。