Skip to content

Commit

Permalink
Add serializer
Browse files Browse the repository at this point in the history
  • Loading branch information
damccorm committed Dec 19, 2024
1 parent 90c13e9 commit a9eed7d
Showing 1 changed file with 3 additions and 1 deletion.
Original file line number Diff line number Diff line change
Expand Up @@ -109,6 +109,7 @@ public PTransform<PCollection<byte[]>, KinesisIO.Write.Result> buildExternal(
AwsBasicCredentials creds =
AwsBasicCredentials.create(configuration.awsAccessKey, configuration.awsSecretKey);
StaticCredentialsProvider provider = StaticCredentialsProvider.create(creds);
SerializableFunction<byte[], byte[]> serializer = v -> v;
KinesisIO.Write<byte[]> writeTransform =
KinesisIO.<byte[]>write()
.withStreamName(configuration.streamName)
Expand All @@ -118,7 +119,8 @@ public PTransform<PCollection<byte[]>, KinesisIO.Write.Result> buildExternal(
.region(configuration.region)
.endpoint(configuration.serviceEndpoint)
.build())
.withPartitioner(p -> configuration.partitionKey);
.withPartitioner(p -> configuration.partitionKey)
.withSerializer(serializer);

return writeTransform;
}
Expand Down

0 comments on commit a9eed7d

Please sign in to comment.