Skip to content
New issue

Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.

By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.

Already on GitHub? Sign in to your account

feat(connector): add DynamoDB sink #16670

Merged
merged 13 commits into from
Jun 4, 2024
79 changes: 56 additions & 23 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

1 change: 1 addition & 0 deletions src/connector/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,7 @@ await-tree = { workspace = true }
aws-config = { workspace = true }
aws-credential-types = { workspace = true }
aws-msk-iam-sasl-signer = "1.0.0"
aws-sdk-dynamodb = "1.23.0"
aws-sdk-kinesis = { workspace = true }
aws-sdk-s3 = { workspace = true }
aws-smithy-http = { workspace = true }
Expand Down
43 changes: 43 additions & 0 deletions src/connector/src/connector_common/common.rs
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@ use std::time::Duration;
use anyhow::{anyhow, Context};
use async_nats::jetstream::consumer::DeliverPolicy;
use async_nats::jetstream::{self};
use aws_sdk_dynamodb::client::Client as DynamoDBClient;
use aws_sdk_kinesis::Client as KinesisClient;
use pulsar::authentication::oauth2::{OAuth2Authentication, OAuth2Params};
use pulsar::{Authentication, Pulsar, TokioExecutor};
Expand Down Expand Up @@ -704,6 +705,48 @@ impl NatsCommon {
}
}

#[derive(Deserialize, Debug, Clone, WithOptions)]
pub struct DynamoDbCommon {
#[serde(rename = "table", alias = "dynamodb.table")]
pub table: String,
// #[serde(rename = "primary_key", alias = "dynamodb.primary_key")]
// pub primary_key: String,
// #[serde(rename = "sort_key", alias = "dynamodb.sort_key")]
// pub sort_key: Option<String>,
#[serde(rename = "aws.region")]
pub stream_region: String,
#[serde(rename = "aws.endpoint")]
pub endpoint: Option<String>,
#[serde(rename = "aws.credentials.access_key_id")]
pub credentials_access_key: Option<String>,
#[serde(rename = "aws.credentials.secret_access_key")]
pub credentials_secret_access_key: Option<String>,
#[serde(rename = "aws.credentials.session_token")]
pub session_token: Option<String>,
#[serde(rename = "aws.credentials.role.arn")]
pub assume_role_arn: Option<String>,
#[serde(rename = "aws.credentials.role.external_id")]
pub assume_role_external_id: Option<String>,
xxhZs marked this conversation as resolved.
Show resolved Hide resolved
}

impl DynamoDbCommon {
pub(crate) async fn build_client(&self) -> ConnectorResult<DynamoDBClient> {
let config = AwsAuthProps {
region: Some(self.stream_region.clone()),
endpoint: self.endpoint.clone(),
access_key: self.credentials_access_key.clone(),
secret_key: self.credentials_secret_access_key.clone(),
session_token: self.session_token.clone(),
arn: self.assume_role_arn.clone(),
external_id: self.assume_role_external_id.clone(),
profile: Default::default(),
};
let aws_config = config.build_config().await?;

Ok(DynamoDBClient::new(&aws_config))
}
}

pub(crate) fn load_certs(
certificates: &str,
) -> ConnectorResult<Vec<rustls_pki_types::CertificateDer<'static>>> {
Expand Down
Loading
Loading