apache/pinot

Support consumer de-aggregation in Kinesis

開放

#10,152 建立於 2023年1月19日

 (0 則留言) (0 個反應) (0 位負責人)Java (1,234 個分叉)batch import
enhancementhelp wantedkinesis

倉庫指標

星標
 (4,937 顆星)
PR 合併指標
 (PR 指標待抓取)

描述

Kinesis supports producing aggregated (aka batched) record. Thus, kinesis consumer also has support for de-aggregating the records. Refer - https://docs.aws.amazon.com/streams/latest/dev/kinesis-kpl-consumer-deaggregation.html

Kinesis provides an AggregatorUtil (https://github.com/awslabs/amazon-kinesis-client/blob/master/amazon-kinesis-client/src/main/java/software/amazon/kinesis/retrieval/AggregatorUtil.java) that can be used in the KinesisConsumer implementation. An example usage of this util can be found in the beam repo (https://github.com/apache/beam/blob/master/sdks/java/io/amazon-web-services2/src/main/java/org/apache/beam/sdk/io/aws2/kinesis/SimplifiedKinesisClient.java#L269)

Even though the records in the batch have the same sequence number, we can append the sub-sequence number to Pinot kinesis' StreamMessageOffset. Changes should be fairly trivial to do this.

Labels: enhancement , kinesis, help-wanted

貢獻者指南