From 936b07faca3e5bc7d7804b0ce7f0a7bb7ef914e5 Mon Sep 17 00:00:00 2001 From: Gabor Barna Date: Wed, 14 Oct 2020 16:06:49 +0200 Subject: [PATCH] kinesis source timestamp (in milliseconds) starting position --- README.md | 2 +- .../org/apache/spark/sql/kinesis/KinesisSourceProvider.scala | 2 ++ 2 files changed, 3 insertions(+), 1 deletion(-) diff --git a/README.md b/README.md index 4a845f0..472a631 100644 --- a/README.md +++ b/README.md @@ -123,7 +123,7 @@ Refering $SPARK_HOME to the Spark installation directory. | awsSTSRoleARN | - | AWS STS Role ARN for Kinesis describe, read record operations | | awsSTSSessionName | - | AWS STS Session name for Kinesis describe, read record operations | | awsUseInstanceProfile | true | Use Instance Profile Credentials if none of credentials provided | -| startingPosition | LATEST | Starting Position in Kinesis to fetch data from. Possible values are "latest", "trim_horizon", "earliest" (alias for trim_horizon), or JSON serialized map shardId->KinesisPosition | +| startingPosition | LATEST | Starting Position in Kinesis to fetch data from. Possible values are "latest", "trim_horizon", "earliest" (alias for trim_horizon), timestamp in milliseconds, or JSON serialized map shardId->KinesisPosition | | failondataloss| true | fail the streaming job if any active shard is missing or expired | kinesis.executor.maxFetchTimeInMs | 1000 | Maximum time spent in executor to fetch record from Kinesis per Shard | | kinesis.executor.maxFetchRecordsPerShard | 100000 | Maximum Number of records to fetch per shard | diff --git a/src/main/scala/org/apache/spark/sql/kinesis/KinesisSourceProvider.scala b/src/main/scala/org/apache/spark/sql/kinesis/KinesisSourceProvider.scala index 9912352..c0ac3c9 100644 --- a/src/main/scala/org/apache/spark/sql/kinesis/KinesisSourceProvider.scala +++ b/src/main/scala/org/apache/spark/sql/kinesis/KinesisSourceProvider.scala @@ -189,6 +189,8 @@ private[kinesis] object KinesisSourceProvider extends Logging { InitialKinesisPosition.fromPredefPosition(new TrimHorizon) case Some(position) if position.toLowerCase(Locale.ROOT) == "earliest" => InitialKinesisPosition.fromPredefPosition(new TrimHorizon) + case Some(timestamp) if timestamp.forall(_.isDigit) => + InitialKinesisPosition.fromPredefPosition(new AtTimeStamp(timestamp)) case Some(json) => InitialKinesisPosition.fromCheckpointJson(json, new AtTimeStamp(CURRENT_TIMESTAMP)) case None => InitialKinesisPosition.fromPredefPosition(new AtTimeStamp(CURRENT_TIMESTAMP))