DataHub成本节省攻略

更新时间:
复制 MD 格式

本文介绍了DataHub新版本带来的改动,介绍batch的原理和实现,以及使用batch后所带来的性能的提升和费用的减少。切换为batch后,对于DataHub而言,服务端的资源消耗会明显降低,同时,性能会明显提升,使用费用也会大幅降低。

升级内容

支持zstd压缩

DataHub在新版本中对zstd压缩算法做了支持,相较于DataHub支持的lz4deflate压缩算法,效果卓越。

说明

zstd是一种高性能压缩算法,由Facebook开发,于2016年开源,zstd在压缩速度和压缩比两方面都有不俗的表现,非常契合DataHub的使用场景。

序列化改造

DataHub引入了batch序列化,batch序列化本质上就是DataHub对数据传输中数据的定义的一种组织方式,batch并不是特指某种序列化的方式,而是对序列化的数据做了二次封装。例如:一次发送100条数据,将100条数据序列化后得到一个buffer,给这个buffer选择一个压缩算法得到压缩后的buffer,这个时候给这个压缩后的buffer添加一个header记录这个buffer大小、数据条数、压缩算法、crc等信息,从而获得一条完整batch buffer。

解决的问题:

  • 可以有效避免业务层面的脏数据。

  • 减少服务端CPU开销,提高数据处理性能。

  • 较少同时读写延迟。

image
说明

batchbuffer发送到服务端后,因客户端已经做了充分的数据有效性的校验,所以服务端只需检验数据中的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的收费项主要是这两个维度为存储与流量,其他主要是为了防止滥用而设置的惩罚性质的收费项。因此,以上测试结果从存储流量两个角度进行分析。

  • 从存储成本上来看,DataHubprotobuf序列化是没有存储压缩的(只是HTTP传输环节压缩),如果替换为batch+zstd,那么存储会由11506KB降为1112KB,也就是说,存储成本下降幅度达到约90%。

  • 从流量成本上来看,DataHubprotobuf+lz4后的大小为3050KB,batch+zstd的大小为1112KB,也就是说,流量成本会降低约60%。

重要

以上为样本数据测试结果,不同数据测试效果有差异,请您根据您业务实际情况进行测试。

使用batch

说明

注意事项

  • batch写入最大的优势需要充分攒批,如果客户端无法攒批,或者攒批的数据较少,可能带来的效果提升并不显著。

  • 为了让用户迁移更加方便,DataHub在各种读写方式之间做了兼容,保证用户中间状态可以更平滑地过渡,即batch写入依旧可以使用原方式读取,原方式写入依旧可以使用batch读取。因此,在写入端更新为batch写入之后,最好消费端也更新为batch,写入和消费不对应反而会降低性能

重要

DataHub新版本已经取消只有打开多versionTopic 才可以使用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。