本文介绍了DataHub新版本带来的改动,介绍batch的原理和实现,以及使用batch后所带来的性能的提升和费用的减少。切换为batch后,对于DataHub而言,服务端的资源消耗会明显降低,同时,性能会明显提升,使用费用也会大幅降低。
升级内容
支持zstd压缩
DataHub在新版本中对zstd压缩算法做了支持,相较于DataHub支持的lz4和deflate压缩算法,效果卓越。
zstd是一种高性能压缩算法,由Facebook开发,于2016年开源,zstd在压缩速度和压缩比两方面都有不俗的表现,非常契合DataHub的使用场景。
序列化改造
DataHub引入了batch序列化,batch序列化本质上就是DataHub对数据传输中数据的定义的一种组织方式,batch并不是特指某种序列化的方式,而是对序列化的数据做了二次封装。例如:一次发送100条数据,将100条数据序列化后得到一个buffer,给这个buffer选择一个压缩算法得到压缩后的buffer,这个时候给这个压缩后的buffer添加一个header记录这个buffer大小、数据条数、压缩算法、crc等信息,从而获得一条完整batch buffer。
解决的问题:
可以有效避免业务层面的脏数据。
减少服务端CPU开销,提高数据处理性能。
较少同时读写延迟。
batch的buffer发送到服务端后,因客户端已经做了充分的数据有效性的校验,所以服务端只需检验数据中的crc确认为有效buffer后,便可以直接落盘,省去了序列化、反序列化、加解压以及校验的操作,服务端性能提升超过80%,因为是多条数据一起压缩的,所以压缩率也提高了,存储成本也降低了。
费用对比
验证batch所带来的实际效果,进行以下测试,假设场景如下:
测试数据为广告投放相关的数据,大约200列,数据中null比例大约20%~30%。
1000条数据一个batch。
batch内部的序列化使用的是avro。
lz4是之前版本默认的压缩算法,压缩使用zstd来替代lz4。
测试结果如下:
数据源大小(Byte) | lz4压缩(Byte) | zstd压缩(Byte) | |
protobuf序列化 | 11,506,677 | 3,050,640 | 1,158,868 |
batch序列化 | 11,154,596 | 2,931,729 | 1,112,693 |
Datahub的收费项主要是这两个维度为存储与流量,其他主要是为了防止滥用而设置的惩罚性质的收费项。因此,以上测试结果从存储与流量两个角度进行分析。
从存储成本上来看,DataHub的protobuf序列化是没有存储压缩的(只是HTTP传输环节压缩),如果替换为batch+zstd,那么存储会由11506KB降为1112KB,也就是说,存储成本下降幅度达到约90%。
从流量成本上来看,DataHub的protobuf+lz4后的大小为3050KB,batch+zstd的大小为1112KB,也就是说,流量成本会降低约60%。
以上为样本数据测试结果,不同数据测试效果有差异,请您根据您业务实际情况进行测试。
使用batch
注意事项
batch写入最大的优势需要充分攒批,如果客户端无法攒批,或者攒批的数据较少,可能带来的效果提升并不显著。
为了让用户迁移更加方便,DataHub在各种读写方式之间做了兼容,保证用户中间状态可以更平滑地过渡,即batch写入依旧可以使用原方式读取,原方式写入依旧可以使用batch读取。因此,在写入端更新为batch写入之后,最好消费端也更新为batch,写入和消费不对应反而会降低性能。
DataHub新版本已经取消只有打开多version的Topic 才可以使用batch协议,无需打开Topic 多version开关就可以使用batch协议,但需要用户更新到最新版的客户端,目前支持的客户端以及所需版本如下所示:
语言 | 支持状态 | 版本要求 |
Java SDK | 支持 | 1.5.1 及以上 |
Go SDk | 支持 | 1.1.0 及以上 |
Python | 暂不支持 | - |
C++ | 暂不支持 | - |
Java 示例:
在 1.5.1 及以后的版本中,默认序列化协议已经设置为 Batch,用户无需做任何特殊配置,只需要正常的写入即可,同步和异步写入均可以。
Maven依赖
<dependency>
<groupId>com.aliyun.datahub</groupId>
<artifactId>datahub-client-library</artifactId>
<version>1.5.1</version>
</dependency>代码示例:
public static void main(String[] args) throws InterruptedException {
// 通过环境变量获取AK信息
EnvironmentVariableCredentialProvider provider = EnvironmentVariableCredentialProvider.create();
String endpoint ="https://dh-cn-hangzhou.aliyuncs.com";
String projectName = "test_project";
String topicName = "test_topic";
// 初始化Producer,这里直接使用默认配置
ProducerConfig config = new ProducerConfig(endpoint, provider);
DatahubProducer producer = new DatahubProducer(projectName, topicName, config);
RecordSchema schema = producer.getTopicSchema();
// 如果开启了多version schema,这里也可以获取指定version的schema
// RecordSchema schema = producer.getTopicSchema(3);
// 对于异步写入,可以根据需要来选择是否注册回调函数
WriteCallback callback = new WriteCallback() {
@Override
public void onSuccess(String shardId, List<RecordEntry> records, long elapsedTimeMs, long sendTimeMs) {
System.out.println("write success");
}
@Override
public void onFailure(String shardId, List<RecordEntry> records, long elapsedTimeMs, DatahubClientException e) {
System.out.println("write failed");
}
};
for (int i = 0; i < 10000; ++i) {
try {
// generate data by schema
TupleRecordData data = new TupleRecordData(schema);
data.setField("field1", "hello");
data.setField("field2", 1234);
RecordEntry recordEntry = new RecordEntry();
recordEntry.setRecordData(data);
producer.sendAsync(recordEntry, callback);
// 如果不需要关心数据是否发送成功,那么就不需要注册回调,直接发送
// producer.sendAsync(recordEntry, null);
} catch (DatahubClientException e) {
// TODO 处理异常,一般是不可重试错误或者超过重试次数;
Thread.sleep(1000);
}
}
// 保证退出前,数据全部被发送完
producer.flush(true);
producer.close();
}Go 示例:
go.mod依赖
require (
github.com/aliyun/aliyun-datahub-sdk-go v1.1.0
)代码示例:
在 1.1.0 及以后的版本中,默认序列化协议已经设置为 Batch,无需用户做任何特殊配置,只需要正常的写入即可,同步和异步写入均可以。
func handleSuccessRun(producer datahub.AsyncProducer) {
for suc := range producer.Successes() {
// handle request success
fmt.Printf("shard:%s, rid:%s, records:%d, latency:%v\n",
suc.ShardId, suc.RequestId, len(suc.Records), suc.Latency)
}
}
func handleFailedRun(producer datahub.AsyncProducer) {
// handle request failed
for err := range producer.Errors() {
fmt.Printf("shard:%s, records:%d, latency:%v, error:%v\n",
err.ShardId, len(err.Records), err.Latency, err.Err)
}
}
func main() {
cfg := datahub.NewProducerConfig()
// 通过环境变量获取AK信息
credential, err := credentials.NewCredential(nil)
if err != nil {
fmt.Println(err)
// TODO: handle error
}
cfg.Account = datahub.NewCredentialAccount(credential)
cfg.Endpoint = "https://dh-cn-hangzhou.aliyuncs.com"
cfg.Project = "test_project"
cfg.Topic = "test_topic"
producer := datahub.NewAsyncProducer(cfg)
err = producer.Init()
if err != nil {
// TODO: handle error
fmt.Println(err)
}
schema, err := producer.GetSchema()
if err != nil {
// TODO: handle error
fmt.Println(err)
}
// 处理success channle
go handleSuccessRun(producer)
// 处理error channle
go handleFailedRun(producer)
// 循环生成数据
for i := 0; i < 1000; i++ {
record := datahub.NewTupleRecord(schema)
// 每次设置值都需要检查设置是否成功
err = record.SetValueByName("f1", "val1")
if err != nil {
fmt.Println(err)
return
}
err = record.SetValueByName("f2", 1234)
if err != nil {
fmt.Println(err)
return
}
producer.Input() <- record
}
err = producer.Close()
if err != nil {
fmt.Println(err)
}
}
后续支持
如果您在使用中有任何问题或者疑问,欢迎提工单咨询或加入用户群咨询,群号:33517130。