apache/pinot

Support consumer de-aggregation in Kinesis

Ouverte

#10 152 ouverte le 19 janv. 2023

 (0 commentaire) (0 réaction) (0 personne assignée)Java (1 234 forks)batch import
enhancementhelp wantedkinesis

Métriques du dépôt

Stars
 (4 937 étoiles)
Métriques de merge PR
 (Métriques PR en attente)

Description

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

Guide contributeur