Skip to content

Commit

Permalink
fix(connector): make kinesis work again (#6317)
Browse files Browse the repository at this point in the history
* add kinesis props

* fix comments

Signed-off-by: tabVersion <[email protected]>

Signed-off-by: tabVersion <[email protected]>
Co-authored-by: mergify[bot] <37929162+mergify[bot]@users.noreply.github.com>
  • Loading branch information
tabVersion and mergify[bot] authored Nov 12, 2022
1 parent 6727507 commit 001387d
Show file tree
Hide file tree
Showing 2 changed files with 37 additions and 1 deletion.
9 changes: 9 additions & 0 deletions src/connector/src/source/kinesis/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -52,4 +52,13 @@ pub struct KinesisProperties {
alias = "kinesis.assumerole.external_id"
)]
pub assume_role_external_id: Option<String>,

#[serde(rename = "scan.startup.mode", alias = "kinesis.scan.startup.mode")]
// accepted values: "latest", "earliest", "sequence_number"
pub scan_startup_mode: Option<String>,
#[serde(
rename = "scan.startup.sequence_number",
alias = "kinesis.scan.startup.sequence_number"
)]
pub seq_offset: Option<String>,
}
29 changes: 28 additions & 1 deletion src/connector/src/source/kinesis/source/reader.rs
Original file line number Diff line number Diff line change
Expand Up @@ -54,15 +54,40 @@ impl SplitReader for KinesisSplitReader {
SplitImpl::Kinesis(ks) => ks,
split => return Err(anyhow!("expect KinesisSplit, got {:?}", split)),
};

let start_position = match &split.start_position {
KinesisOffset::None => match &properties.scan_startup_mode {
None => KinesisOffset::Earliest,
Some(mode) => match mode.as_str() {
"earliest" => KinesisOffset::Earliest,
"latest" => KinesisOffset::Latest,
"sequence_number" => {
if let Some(seq) = &properties.seq_offset {
KinesisOffset::SequenceNumber(seq.clone())
} else {
return Err(anyhow!("scan_startup_sequence_number is required"));
}
}
_ => {
return Err(anyhow!(
"invalid scan_startup_mode, accept earliest/latest/sequence_number"
))
}
},
},
start_position => start_position.to_owned(),
};

let stream_name = properties.stream_name.clone();
let client = build_client(properties).await?;

Ok(Self {
client,
stream_name,
shard_id: split.shard_id,
shard_iter: None,
latest_offset: None,
start_position: split.start_position,
start_position,
end_position: split.end_position,
})
}
Expand Down Expand Up @@ -166,6 +191,8 @@ mod tests {
endpoint: None,
session_token: None,
assume_role_external_id: None,
scan_startup_mode: None,
seq_offset: None,
};

let mut trim_horizen_reader = KinesisSplitReader::new(
Expand Down

0 comments on commit 001387d

Please sign in to comment.