Skip to content

Commit

Permalink
CAMEL-17669 - AWS Kinesis: Make the consumer set only the raw data on…
Browse files Browse the repository at this point in the history
… the body, and enrich the exchange with more header
  • Loading branch information
oscerd committed Mar 11, 2022
1 parent b931ab7 commit e5cf3d4
Show file tree
Hide file tree
Showing 2 changed files with 9 additions and 2 deletions.
Original file line number Diff line number Diff line change
Expand Up @@ -191,7 +191,7 @@ private Queue<Exchange> createExchanges(List<Record> records) {

protected Exchange createExchange(Record record) {
Exchange exchange = createExchange(true);
exchange.getIn().setBody(record);
exchange.getIn().setBody(record.data().asInputStream());
exchange.getIn().setHeader(Kinesis2Constants.APPROX_ARRIVAL_TIME, record.approximateArrivalTimestamp());
exchange.getIn().setHeader(Kinesis2Constants.PARTITION_KEY, record.partitionKey());
exchange.getIn().setHeader(Kinesis2Constants.SEQUENCE_NUMBER, record.sequenceNumber());
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@
*/
package org.apache.camel.component.aws2.kinesis;

import java.nio.charset.Charset;
import java.util.ArrayList;

import org.apache.camel.AsyncProcessor;
Expand All @@ -27,6 +28,7 @@
import org.mockito.ArgumentCaptor;
import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension;
import software.amazon.awssdk.core.SdkBytes;
import software.amazon.awssdk.services.kinesis.KinesisClient;
import software.amazon.awssdk.services.kinesis.model.DescribeStreamRequest;
import software.amazon.awssdk.services.kinesis.model.DescribeStreamResponse;
Expand Down Expand Up @@ -78,7 +80,12 @@ public void setup() {

when(kinesisClient.getRecords(any(GetRecordsRequest.class))).thenReturn(GetRecordsResponse.builder()
.nextShardIterator("nextShardIterator")
.records(Record.builder().sequenceNumber("1").build(), Record.builder().sequenceNumber("2").build()).build());
.records(
Record.builder().sequenceNumber("1").data(SdkBytes.fromString("Hello", Charset.defaultCharset()))
.build(),
Record.builder().sequenceNumber("2").data(SdkBytes.fromString("Hello", Charset.defaultCharset()))
.build())
.build());
when(kinesisClient.describeStream(any(DescribeStreamRequest.class)))
.thenReturn(DescribeStreamResponse.builder()
.streamDescription(StreamDescription.builder().shards(shardList).build()).build());
Expand Down

0 comments on commit e5cf3d4

Please sign in to comment.