diff --git a/Cargo.lock b/Cargo.lock index c6aa065741..a1fe867ee1 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1336,7 +1336,7 @@ dependencies = [ "base64", "http 1.4.2", "log", - "rustls", + "rustls 0.23.42", "serde", "serde_json", "url", @@ -1687,17 +1687,23 @@ dependencies = [ "aws-smithy-async", "aws-smithy-runtime-api", "aws-smithy-types", + "h2 0.3.27", "h2 0.4.15", + "http 0.2.12", "http 1.4.2", - "hyper", - "hyper-rustls", + "http-body 0.4.6", + "hyper 0.14.32", + "hyper 1.11.0", + "hyper-rustls 0.24.2", + "hyper-rustls 0.27.9", "hyper-util", "pin-project-lite", - "rustls", + "rustls 0.21.12", + "rustls 0.23.42", "rustls-native-certs", "rustls-pki-types", "tokio", - "tokio-rustls", + "tokio-rustls 0.26.4", "tower", "tracing", ] @@ -1865,7 +1871,7 @@ dependencies = [ "http 1.4.2", "http-body 1.1.0", "http-body-util", - "hyper", + "hyper 1.11.0", "hyper-util", "itoa", "matchit", @@ -1927,13 +1933,13 @@ dependencies = [ "fs-err", "http 1.4.2", "http-body 1.1.0", - "hyper", + "hyper 1.11.0", "hyper-util", "pin-project-lite", - "rustls", + "rustls 0.23.42", "rustls-pki-types", "tokio", - "tokio-rustls", + "tokio-rustls 0.26.4", "tower-service", ] @@ -2305,16 +2311,16 @@ dependencies = [ "home", "http 1.4.2", "http-body-util", - "hyper", + "hyper 1.11.0", "hyper-named-pipe", - "hyper-rustls", + "hyper-rustls 0.27.9", "hyper-util", "hyperlocal", "log", "num", "pin-project-lite", "rand 0.9.5", - "rustls", + "rustls 0.23.42", "rustls-native-certs", "rustls-pki-types", "serde", @@ -3142,7 +3148,7 @@ dependencies = [ "libc", "quinn-proto", "rustc-hash", - "rustls", + "rustls 0.23.42", "rustls-platform-verifier", "synchrony", "thiserror 2.0.19", @@ -3189,7 +3195,7 @@ dependencies = [ "compio-io", "futures-rustls", "futures-util", - "rustls", + "rustls 0.23.42", ] [[package]] @@ -3796,7 +3802,7 @@ dependencies = [ "futures-util", "http 1.4.2", "http-body-util", - "hyper", + "hyper 1.11.0", "hyper-util", "mime", "rustls-platform-verifier", @@ -3821,7 +3827,7 @@ dependencies = [ "compio-log", "cyper-core", "futures-util", - "hyper", + "hyper 1.11.0", "hyper-util", "send_wrapper", "socket2 0.6.5", @@ -3838,7 +3844,7 @@ checksum = "03c8847069e286c64987119637d5f08cdb71e12e85be0294fed649dc8007d32e" dependencies = [ "compio", "futures-util", - "hyper", + "hyper 1.11.0", "send_wrapper", ] @@ -5278,7 +5284,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "a8f2f12607f92c69b12ed746fabf9ca4f5c482cba46679c1a75b874ed7c26adb" dependencies = [ "futures-io", - "rustls", + "rustls 0.23.42", "rustls-pki-types", ] @@ -6317,6 +6323,30 @@ dependencies = [ "typenum", ] +[[package]] +name = "hyper" +version = "0.14.32" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "41dfc780fdec9373c01bae43289ea34c972e40ee3c9f6b3c8801a35f35586ce7" +dependencies = [ + "bytes", + "futures-channel", + "futures-core", + "futures-util", + "h2 0.3.27", + "http 0.2.12", + "http-body 0.4.6", + "httparse", + "httpdate", + "itoa", + "pin-project-lite", + "socket2 0.5.10", + "tokio", + "tower-service", + "tracing", + "want", +] + [[package]] name = "hyper" version = "1.11.0" @@ -6346,13 +6376,28 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "fab3637d6b04a8037af8a266fdf6cf92ea957e8c53981a2bf6136572531025bf" dependencies = [ "hex", - "hyper", + "hyper 1.11.0", "hyper-util", "pin-project-lite", "tokio", "tower-service", ] +[[package]] +name = "hyper-rustls" +version = "0.24.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ec3efd23720e2049821a693cbc7e65ea87c72f1c58ff2f9522ff332b1491e590" +dependencies = [ + "futures-util", + "http 0.2.12", + "hyper 0.14.32", + "log", + "rustls 0.21.12", + "tokio", + "tokio-rustls 0.24.1", +] + [[package]] name = "hyper-rustls" version = "0.27.9" @@ -6360,13 +6405,13 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "33ca68d021ef39cf6463ab54c1d0f5daf03377b70561305bb89a8f83aab66e0f" dependencies = [ "http 1.4.2", - "hyper", + "hyper 1.11.0", "hyper-util", "log", - "rustls", + "rustls 0.23.42", "rustls-native-certs", "tokio", - "tokio-rustls", + "tokio-rustls 0.26.4", "tower-service", "webpki-roots 1.0.9", ] @@ -6377,7 +6422,7 @@ version = "0.5.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "2b90d566bffbce6a75bd8b09a05aa8c2cb1fabb6cb348f8840c9e4c90a0d83b0" dependencies = [ - "hyper", + "hyper 1.11.0", "hyper-util", "pin-project-lite", "tokio", @@ -6396,7 +6441,7 @@ dependencies = [ "futures-util", "http 1.4.2", "http-body 1.1.0", - "hyper", + "hyper 1.11.0", "ipnet", "libc", "percent-encoding", @@ -6417,7 +6462,7 @@ checksum = "986c5ce3b994526b3cd75578e62554abd09f0899d6206de48b3e96ab34ccc8c7" dependencies = [ "hex", "http-body-util", - "hyper", + "hyper 1.11.0", "hyper-util", "pin-project-lite", "tokio", @@ -6679,12 +6724,12 @@ dependencies = [ "reqwest-middleware", "reqwest-retry", "reqwest-tracing", - "rustls", + "rustls 0.23.42", "secrecy", "serde", "serde_json", "tokio", - "tokio-rustls", + "tokio-rustls 0.26.4", "tokio-tungstenite", "tracing", "trait-variant", @@ -6991,6 +7036,22 @@ dependencies = [ "wiremock", ] +[[package]] +name = "iggy_connector_dynamodb_sink" +version = "0.5.0-edge.4" +dependencies = [ + "async-trait", + "aws-config", + "aws-sdk-dynamodb", + "humantime", + "iggy_connector_sdk", + "secrecy", + "serde", + "simd-json", + "tokio", + "tracing", +] + [[package]] name = "iggy_connector_elasticsearch_sink" version = "0.5.0-edge.4" @@ -7526,6 +7587,8 @@ dependencies = [ "arrow 57.3.1", "assert_cmd", "async-trait", + "aws-config", + "aws-sdk-dynamodb", "base64", "bon", "bytemuck", @@ -8580,7 +8643,7 @@ dependencies = [ "rand 0.10.2", "rcgen", "ring", - "rustls", + "rustls 0.23.42", "rustls-pemfile", "scopeguard", "server_common", @@ -8796,7 +8859,7 @@ dependencies = [ "percent-encoding", "rand 0.9.5", "rustc_version_runtime", - "rustls", + "rustls 0.23.42", "serde", "serde_bytes", "serde_with", @@ -8808,7 +8871,7 @@ dependencies = [ "take_mut", "thiserror 2.0.19", "tokio", - "tokio-rustls", + "tokio-rustls 0.26.4", "tokio-util", "typed-builder 0.22.0", "uuid", @@ -9211,7 +9274,7 @@ dependencies = [ "http-body-util", "httparse", "humantime", - "hyper", + "hyper 1.11.0", "itertools 0.14.0", "md-5 0.10.6", "parking_lot", @@ -9253,8 +9316,8 @@ dependencies = [ "http-body 1.1.0", "http-body-util", "http-serde", - "hyper", - "hyper-rustls", + "hyper 1.11.0", + "hyper-rustls 0.27.9", "hyper-timeout", "hyper-util", "jsonwebtoken", @@ -9874,7 +9937,7 @@ dependencies = [ "stringprep", "thiserror 2.0.19", "tokio", - "tokio-rustls", + "tokio-rustls 0.26.4", "tokio-util", "x509-certificate", ] @@ -10565,7 +10628,7 @@ dependencies = [ "quinn-proto", "quinn-udp", "rustc-hash", - "rustls", + "rustls 0.23.42", "socket2 0.6.5", "thiserror 2.0.19", "tokio", @@ -10588,7 +10651,7 @@ dependencies = [ "rand_pcg", "ring", "rustc-hash", - "rustls", + "rustls 0.23.42", "rustls-pki-types", "rustls-platform-verifier", "slab", @@ -10979,15 +11042,15 @@ dependencies = [ "http 1.4.2", "http-body 1.1.0", "http-body-util", - "hyper", - "hyper-rustls", + "hyper 1.11.0", + "hyper-rustls 0.27.9", "hyper-util", "js-sys", "log", "percent-encoding", "pin-project-lite", "quinn", - "rustls", + "rustls 0.23.42", "rustls-native-certs", "rustls-pki-types", "serde", @@ -10995,7 +11058,7 @@ dependencies = [ "serde_urlencoded", "sync_wrapper", "tokio", - "tokio-rustls", + "tokio-rustls 0.26.4", "tokio-util", "tower", "tower-http 0.6.11", @@ -11024,8 +11087,8 @@ dependencies = [ "http 1.4.2", "http-body 1.1.0", "http-body-util", - "hyper", - "hyper-rustls", + "hyper 1.11.0", + "hyper-rustls 0.27.9", "hyper-util", "js-sys", "log", @@ -11033,7 +11096,7 @@ dependencies = [ "percent-encoding", "pin-project-lite", "quinn", - "rustls", + "rustls 0.23.42", "rustls-pki-types", "rustls-platform-verifier", "serde", @@ -11041,7 +11104,7 @@ dependencies = [ "serde_urlencoded", "sync_wrapper", "tokio", - "tokio-rustls", + "tokio-rustls 0.26.4", "tokio-util", "tower", "tower-http 0.6.11", @@ -11079,7 +11142,7 @@ dependencies = [ "futures", "getrandom 0.2.17", "http 1.4.2", - "hyper", + "hyper 1.11.0", "reqwest 0.13.4", "reqwest-middleware", "retry-policies", @@ -11478,6 +11541,18 @@ dependencies = [ "windows-sys 0.61.2", ] +[[package]] +name = "rustls" +version = "0.21.12" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3f56a14d1f48b391359b22f731fd4bd7e43c97f3c50eee276f3aa09c94784d3e" +dependencies = [ + "log", + "ring", + "rustls-webpki 0.101.7", + "sct", +] + [[package]] name = "rustls" version = "0.23.42" @@ -11489,7 +11564,7 @@ dependencies = [ "once_cell", "ring", "rustls-pki-types", - "rustls-webpki", + "rustls-webpki 0.103.13", "subtle", "zeroize", ] @@ -11536,10 +11611,10 @@ dependencies = [ "jni", "log", "once_cell", - "rustls", + "rustls 0.23.42", "rustls-native-certs", "rustls-platform-verifier-android", - "rustls-webpki", + "rustls-webpki 0.103.13", "security-framework", "security-framework-sys", "webpki-root-certs", @@ -11552,6 +11627,16 @@ version = "0.1.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "f87165f0995f63a9fbeea62b64d10b4d9d8e78ec6d7d51fb2125fda7bb36788f" +[[package]] +name = "rustls-webpki" +version = "0.101.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8b6275d1ee7a1cd780b64aca7726599a1dbc893b1e64144529e55c3c2f745765" +dependencies = [ + "ring", + "untrusted 0.9.0", +] + [[package]] name = "rustls-webpki" version = "0.103.13" @@ -11681,6 +11766,16 @@ version = "1.2.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "94143f37725109f92c262ed2cf5e59bce7498c01bcc1502d7b9afe439a4e9f49" +[[package]] +name = "sct" +version = "0.7.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "da046153aa2352493d6cb7da4b6e5c0c057d8a1d0a9aa8560baffdd945acd414" +dependencies = [ + "ring", + "untrusted 0.9.0", +] + [[package]] name = "sd-notify" version = "0.5.0" @@ -12063,7 +12158,7 @@ dependencies = [ "hash32 1.0.0", "human-repr", "hwlocality", - "hyper", + "hyper 1.11.0", "hyper-util", "iggy_binary_protocol", "iggy_common", @@ -12088,7 +12183,7 @@ dependencies = [ "rmp-serde", "rolling-file", "rust-embed", - "rustls", + "rustls 0.23.42", "rustls-pemfile", "sd-notify", "secrecy", @@ -12142,7 +12237,7 @@ dependencies = [ "rand 0.10.2", "rcgen", "rolling-file", - "rustls", + "rustls 0.23.42", "send_wrapper", "serde", "serial_test", @@ -12607,7 +12702,7 @@ dependencies = [ "log", "memchr", "percent-encoding", - "rustls", + "rustls 0.23.42", "serde", "serde_json", "sha2 0.10.9", @@ -13471,13 +13566,23 @@ dependencies = [ "whoami", ] +[[package]] +name = "tokio-rustls" +version = "0.24.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c28327cf380ac148141087fbfb9de9d7bd4e84ab5d2c28fbc911d753de8a7081" +dependencies = [ + "rustls 0.21.12", + "tokio", +] + [[package]] name = "tokio-rustls" version = "0.26.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "1729aa945f29d91ba541258c8df89027d5792d85a8841fb65e8bf0f4ede4ef61" dependencies = [ - "rustls", + "rustls 0.23.42", "tokio", ] @@ -13500,10 +13605,10 @@ checksum = "17a073bfed563fa236697a068031408a93cd9522e08abf9933ead3e73411bd71" dependencies = [ "futures-util", "log", - "rustls", + "rustls 0.23.42", "rustls-pki-types", "tokio", - "tokio-rustls", + "tokio-rustls 0.26.4", "tungstenite 0.30.0", "webpki-roots 0.26.11", ] @@ -13645,7 +13750,7 @@ dependencies = [ "http 1.4.2", "http-body 1.1.0", "http-body-util", - "hyper", + "hyper 1.11.0", "hyper-timeout", "hyper-util", "percent-encoding", @@ -13909,7 +14014,7 @@ dependencies = [ "httparse", "log", "rand 0.9.5", - "rustls", + "rustls 0.23.42", "rustls-pki-types", "sha1 0.10.7", "thiserror 2.0.19", @@ -13927,7 +14032,7 @@ dependencies = [ "httparse", "log", "rand 0.10.2", - "rustls", + "rustls 0.23.42", "rustls-pki-types", "sha1 0.11.0", "thiserror 2.0.19", @@ -14211,7 +14316,7 @@ dependencies = [ "flate2", "log", "percent-encoding", - "rustls", + "rustls 0.23.42", "rustls-pki-types", "ureq-proto", "utf8-zero", @@ -15176,7 +15281,7 @@ dependencies = [ "futures", "http 1.4.2", "http-body-util", - "hyper", + "hyper 1.11.0", "hyper-util", "log", "once_cell", diff --git a/Cargo.toml b/Cargo.toml index ea2d476acb..5a9e1596d3 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -36,6 +36,7 @@ members = [ "core/connectors/sinks/clickhouse_sink", "core/connectors/sinks/delta_sink", "core/connectors/sinks/doris_sink", + "core/connectors/sinks/dynamodb_sink", "core/connectors/sinks/elasticsearch_sink", "core/connectors/sinks/http_sink", "core/connectors/sinks/iceberg_sink", @@ -105,6 +106,8 @@ async-channel = "2.5.0" async-dropper = { version = "0.3.1", features = ["tokio", "simple"] } async-trait = "0.1.91" async_zip = { version = "0.0.18", features = ["tokio", "lzma", "bzip2", "xz", "deflate", "zstd"] } +aws-config = { version = "1.9.0", features = ["behavior-version-latest"] } +aws-sdk-dynamodb = "1.117.0" axum = { version = "0.8.9", features = ["macros"] } axum-server = { version = "0.8.0", features = ["tls-rustls"] } base64 = "0.22.1" diff --git a/core/connectors/README.md b/core/connectors/README.md index 890a263b5f..de3c26e80e 100644 --- a/core/connectors/README.md +++ b/core/connectors/README.md @@ -81,6 +81,7 @@ Each sink should have its own, custom configuration, which is passed along with ### Available Sinks - **Doris Sink** - loads JSON messages into Apache Doris tables via the Stream Load HTTP API +- **DynamoDB Sink** - writes messages to Amazon DynamoDB tables - **Elasticsearch Sink** - sends messages to Elasticsearch indices - **Iceberg Sink** - writes data to Apache Iceberg tables via REST catalog - **Meilisearch Sink** - indexes messages in Meilisearch diff --git a/core/connectors/sinks/README.md b/core/connectors/sinks/README.md index 91de1fd084..1e8d6a563e 100644 --- a/core/connectors/sinks/README.md +++ b/core/connectors/sinks/README.md @@ -9,6 +9,7 @@ Sink connectors are responsible for writing data from Iggy streams to external s | Sink | Description | | ---- | ----------- | | **doris_sink** | Loads JSON messages into Apache Doris tables via the Stream Load HTTP API | +| **dynamodb_sink** | Writes messages to Amazon DynamoDB tables with batched, key-deterministic writes | | **elasticsearch_sink** | Sends messages to Elasticsearch indices for full-text search and analytics | | **iceberg_sink** | Writes data to Apache Iceberg tables via REST catalog with S3/GCS/Azure storage | | **influxdb_sink** | Writes messages to InfluxDB as line-protocol points; supports both V2 (org/bucket, Flux) and V3 (db, SQL) | diff --git a/core/connectors/sinks/dynamodb_sink/Cargo.toml b/core/connectors/sinks/dynamodb_sink/Cargo.toml new file mode 100644 index 0000000000..93999f1478 --- /dev/null +++ b/core/connectors/sinks/dynamodb_sink/Cargo.toml @@ -0,0 +1,48 @@ +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you under the Apache License, Version 2.0 (the +# "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, +# software distributed under the License is distributed on an +# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +# KIND, either express or implied. See the License for the +# specific language governing permissions and limitations +# under the License. + +[package] +name = "iggy_connector_dynamodb_sink" +version = "0.5.0-edge.4" +description = "Iggy DynamoDB sink connector for writing stream messages to Amazon DynamoDB tables" +edition = "2024" +license = "Apache-2.0" +keywords = ["iggy", "messaging", "streaming", "dynamodb", "sink"] +categories = ["command-line-utilities", "database", "network-programming"] +homepage = "https://iggy.apache.org" +documentation = "https://iggy.apache.org/docs" +repository = "https://github.com/apache/iggy" +readme = "../../README.md" +publish = false + +[lib] +crate-type = ["cdylib", "lib"] + +[dependencies] +async-trait = { workspace = true } +aws-config = { workspace = true } +aws-sdk-dynamodb = { workspace = true } +humantime = { workspace = true } +iggy_connector_sdk = { workspace = true } +secrecy = { workspace = true } +serde = { workspace = true } +simd-json = { workspace = true } +tokio = { workspace = true } +tracing = { workspace = true } + +[lints] +workspace = true diff --git a/core/connectors/sinks/dynamodb_sink/README.md b/core/connectors/sinks/dynamodb_sink/README.md new file mode 100644 index 0000000000..761094c971 --- /dev/null +++ b/core/connectors/sinks/dynamodb_sink/README.md @@ -0,0 +1,95 @@ +# DynamoDB Sink Connector + +A sink connector that consumes messages from Iggy streams and writes them to an +Amazon DynamoDB table with `BatchWriteItem` through the official AWS SDK. + +## Configuration + +```toml +[plugin_config] +table = "iggy_messages" +region = "us-east-1" +# endpoint = "http://localhost:8000" +# access_key_id = "..." +# secret_access_key = "..." +# session_token = "..." +partition_key_field = "iggy_id" +# sort_key_field = "iggy_offset" +batch_size = 25 +include_metadata = true +include_checksum = true +include_origin_timestamp = true +max_item_size = 409600 +max_retries = 3 +retry_delay = "500ms" +max_retry_delay = "5s" +verbose_logging = false +``` + +- `table`: Target DynamoDB table. The table must already exist. +- `region`: AWS region. Falls back to the default AWS region chain when unset. +- `endpoint`: Custom endpoint URL, for DynamoDB Local or a VPC endpoint. +- `access_key_id` / `secret_access_key` / `session_token`: Static credentials. + Provide both `access_key_id` and `secret_access_key`, or neither. When both + are omitted the connector uses the default AWS credential chain. +- `partition_key_field`: Item attribute used as the table partition key. + Defaults to `iggy_id`. +- `sort_key_field`: Item attribute used as the table sort key. Only set this + when the table has a sort key. +- `batch_size`: Items per `BatchWriteItem` request. Defaults to `25`, which is + also the DynamoDB limit, so larger values are clamped. +- `include_metadata`: Add `iggy_stream`, `iggy_topic`, `iggy_partition_id`, + `iggy_offset`, and `iggy_timestamp` to each item. Defaults to `true`. +- `include_checksum`: Add `iggy_checksum`. Defaults to `true`. +- `include_origin_timestamp`: Add `iggy_origin_timestamp`. Defaults to `true`. +- `max_item_size`: Maximum item size in bytes. Defaults to `409600` (400 KB), + which is also the DynamoDB limit, so larger values are clamped. +- `max_retries`: Retries after the first attempt. Defaults to `3`. +- `retry_delay`: First retry delay as a humantime string. Defaults to `500ms`. +- `max_retry_delay`: Upper bound of the exponential backoff. Defaults to `5s`. +- `verbose_logging`: Log per-batch results at info level. Defaults to `false`. + +## Behavior + +JSON objects are written attribute by attribute, so a message field becomes a +DynamoDB attribute of the matching type. JSON arrays and scalars are nested +under a `payload` attribute, because a DynamoDB item must be a map. Text +payloads go into `payload` as a string. Raw payloads are parsed as JSON when +possible, otherwise they are stored as binary. Protobuf, FlatBuffer, and Avro +payloads are not supported and are skipped with a warning. + +Metadata attributes are written after the payload, so they overwrite payload +fields of the same name. + +## Keys and Idempotency + +`BatchWriteItem` uses `PutRequest`, which overwrites an existing item with the +same primary key. Redelivery of the same message therefore writes the same item +again as long as the key is deterministic. + +When the payload does not carry the configured `partition_key_field`, the +connector injects a key built from the stream, topic, partition, and message ID. +When `sort_key_field` is configured and missing, the message offset is injected. +A payload value always wins over the injected one, so a message that carries the +key field with an empty value, or with a value that is neither a string, a +number, nor binary, is skipped rather than falling back to the injected key, +because DynamoDB would reject the whole batch. + +The key fields are checked against the table on startup. A `partition_key_field` +or `sort_key_field` that does not match the table key schema fails the connector +while it opens, instead of on the first write. + +DynamoDB rejects a request that writes the same key twice, so within one +`consume()` call only the newest item per key is sent. + +## Retries + +Unprocessed items returned by `BatchWriteItem` are retried with exponential +backoff, so a partially throttled batch is not silently dropped. Throttling and +server errors such as `ProvisionedThroughputExceededException`, +`ThrottlingException`, and `InternalServerError` are retried the same way, which +covers both provisioned and on-demand capacity modes. Validation and access +errors are permanent and returned without a retry. + +Items larger than `max_item_size` are logged and skipped instead of failing the +whole batch. diff --git a/core/connectors/sinks/dynamodb_sink/config.toml b/core/connectors/sinks/dynamodb_sink/config.toml new file mode 100644 index 0000000000..22461675e3 --- /dev/null +++ b/core/connectors/sinks/dynamodb_sink/config.toml @@ -0,0 +1,51 @@ +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you under the Apache License, Version 2.0 (the +# "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, +# software distributed under the License is distributed on an +# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +# KIND, either express or implied. See the License for the +# specific language governing permissions and limitations +# under the License. + +type = "sink" +key = "dynamodb" +enabled = true +version = 0 +name = "DynamoDB sink" +path = "../../target/release/libiggy_connector_dynamodb_sink" +verbose = false + +[[streams]] +stream = "test_stream" +topics = ["test_topic"] +schema = "json" +batch_length = 100 +poll_interval = "5ms" +consumer_group = "dynamodb_sink" + +[plugin_config] +table = "iggy_messages" +region = "us-east-1" +# endpoint = "http://localhost:8000" # uncomment for DynamoDB Local +# access_key_id = "..." # omit to use the default AWS credential chain +# secret_access_key = "..." # omit to use the default AWS credential chain +# session_token = "..." # only needed for temporary credentials +partition_key_field = "iggy_id" +# sort_key_field = "iggy_offset" +batch_size = 25 +include_metadata = true +include_checksum = true +include_origin_timestamp = true +max_item_size = 409600 +max_retries = 3 +retry_delay = "500ms" +max_retry_delay = "5s" +verbose_logging = false diff --git a/core/connectors/sinks/dynamodb_sink/src/lib.rs b/core/connectors/sinks/dynamodb_sink/src/lib.rs new file mode 100644 index 0000000000..8e4b22a6f4 --- /dev/null +++ b/core/connectors/sinks/dynamodb_sink/src/lib.rs @@ -0,0 +1,1348 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +use async_trait::async_trait; +use aws_config::BehaviorVersion; +use aws_sdk_dynamodb::Client; +use aws_sdk_dynamodb::config::{Credentials, Region}; +use aws_sdk_dynamodb::error::{ProvideErrorMetadata, SdkError}; +use aws_sdk_dynamodb::primitives::Blob; +use aws_sdk_dynamodb::types::{ + AttributeValue, KeySchemaElement, KeyType, PutRequest, WriteRequest, +}; +use humantime::Duration as HumanDuration; +use iggy_connector_sdk::retry::{exponential_backoff, jitter}; +use iggy_connector_sdk::{ + ConsumedMessage, Error, MessagesMetadata, Payload, Sink, TopicMetadata, sink_connector, +}; +use secrecy::{ExposeSecret, SecretString}; +use serde::Deserialize; +use simd_json::{OwnedValue, StaticNode}; +use std::collections::HashMap; +use std::fmt::Write; +use std::str::FromStr; +use std::sync::atomic::{AtomicU64, Ordering}; +use std::time::Duration; +use tracing::{debug, error, info, warn}; + +sink_connector!(DynamoDbSink); + +/// `BatchWriteItem` rejects any request carrying more than 25 write requests. +const MAX_BATCH_WRITE_ITEMS: usize = 25; +/// DynamoDB rejects items larger than 400 KB. +const MAX_ITEM_SIZE: usize = 400 * 1024; +const DEFAULT_PARTITION_KEY_FIELD: &str = "iggy_id"; +const DEFAULT_MAX_RETRIES: u32 = 3; +const DEFAULT_RETRY_DELAY: &str = "500ms"; +const DEFAULT_MAX_RETRY_DELAY: &str = "5s"; +const PAYLOAD_FIELD: &str = "payload"; +const CREDENTIALS_PROVIDER_NAME: &str = "iggy-dynamodb-sink"; + +#[derive(Debug)] +pub struct DynamoDbSink { + pub id: u32, + client: Option, + config: DynamoDbSinkConfig, + partition_key_field: String, + sort_key_field: Option, + batch_size: usize, + include_metadata: bool, + include_checksum: bool, + include_origin_timestamp: bool, + max_item_size: usize, + max_retries: u32, + retry_delay: Duration, + max_retry_delay: Duration, + verbose: bool, + items_written: AtomicU64, + items_skipped: AtomicU64, + items_deduplicated: AtomicU64, + write_errors: AtomicU64, +} + +/// Only `Deserialize` - nothing serializes a plugin config back out, and the +/// missing impl is what keeps the credentials unserializable. +#[derive(Debug, Clone, Deserialize)] +pub struct DynamoDbSinkConfig { + pub table: String, + pub region: Option, + pub endpoint: Option, + pub access_key_id: Option, + pub secret_access_key: Option, + pub session_token: Option, + pub partition_key_field: Option, + pub sort_key_field: Option, + pub batch_size: Option, + pub include_metadata: Option, + pub include_checksum: Option, + pub include_origin_timestamp: Option, + pub max_item_size: Option, + pub max_retries: Option, + pub retry_delay: Option, + pub max_retry_delay: Option, + pub verbose_logging: Option, +} + +impl DynamoDbSink { + pub fn new(id: u32, config: DynamoDbSinkConfig) -> Self { + let partition_key_field = config + .partition_key_field + .clone() + .unwrap_or_else(|| DEFAULT_PARTITION_KEY_FIELD.to_owned()); + let sort_key_field = config.sort_key_field.clone(); + let batch_size = config + .batch_size + .unwrap_or(MAX_BATCH_WRITE_ITEMS as u32) + .clamp(1, MAX_BATCH_WRITE_ITEMS as u32) as usize; + let include_metadata = config.include_metadata.unwrap_or(true); + let include_checksum = config.include_checksum.unwrap_or(true); + let include_origin_timestamp = config.include_origin_timestamp.unwrap_or(true); + let max_item_size = config + .max_item_size + .unwrap_or(MAX_ITEM_SIZE) + .min(MAX_ITEM_SIZE); + let max_retries = config.max_retries.unwrap_or(DEFAULT_MAX_RETRIES); + let retry_delay = parse_duration(config.retry_delay.as_deref(), DEFAULT_RETRY_DELAY); + let mut max_retry_delay = + parse_duration(config.max_retry_delay.as_deref(), DEFAULT_MAX_RETRY_DELAY); + if max_retry_delay < retry_delay { + warn!( + "DynamoDB sink ID: {id} has max_retry_delay below retry_delay, raising it to the retry delay" + ); + max_retry_delay = retry_delay; + } + let verbose = config.verbose_logging.unwrap_or(false); + + DynamoDbSink { + id, + client: None, + config, + partition_key_field, + sort_key_field, + batch_size, + include_metadata, + include_checksum, + include_origin_timestamp, + max_item_size, + max_retries, + retry_delay, + max_retry_delay, + verbose, + items_written: AtomicU64::new(0), + items_skipped: AtomicU64::new(0), + items_deduplicated: AtomicU64::new(0), + write_errors: AtomicU64::new(0), + } + } +} + +#[async_trait] +impl Sink for DynamoDbSink { + async fn open(&mut self) -> Result<(), Error> { + info!( + "Opening DynamoDB sink connector with ID: {}, table: {}", + self.id, self.config.table + ); + let client = self.build_client().await?; + let description = client + .describe_table() + .table_name(&self.config.table) + .send() + .await + .map_err(|error| { + Error::InitError(format!( + "DynamoDB table '{}' is not reachable, error: {}", + self.config.table, + describe_sdk_error(&error) + )) + })?; + self.validate_key_schema( + description + .table + .and_then(|table| table.key_schema) + .unwrap_or_default(), + )?; + + self.client = Some(client); + info!( + "Opened DynamoDB sink connector with ID: {}, table: {}", + self.id, self.config.table + ); + Ok(()) + } + + async fn consume( + &self, + topic_metadata: &TopicMetadata, + messages_metadata: MessagesMetadata, + messages: Vec, + ) -> Result<(), Error> { + self.write_messages(topic_metadata, &messages_metadata, messages) + .await + } + + async fn close(&mut self) -> Result<(), Error> { + info!("Closing DynamoDB sink connector with ID: {}", self.id); + self.client.take(); + info!( + "Closed DynamoDB sink connector with ID: {}, written: {}, skipped: {}, deduplicated: {}, errors: {}", + self.id, + self.items_written.load(Ordering::Relaxed), + self.items_skipped.load(Ordering::Relaxed), + self.items_deduplicated.load(Ordering::Relaxed), + self.write_errors.load(Ordering::Relaxed) + ); + Ok(()) + } +} + +impl DynamoDbSink { + async fn build_client(&self) -> Result { + if self.config.access_key_id.is_some() != self.config.secret_access_key.is_some() { + return Err(Error::InvalidConfigValue( + "Partially configured credentials. You must provide both access_key_id \ + and secret_access_key, or omit both." + .to_owned(), + )); + } + + let mut loader = aws_config::defaults(BehaviorVersion::latest()); + if let Some(region) = &self.config.region { + loader = loader.region(Region::new(region.clone())); + } + if let Some(endpoint) = &self.config.endpoint { + info!("Using custom DynamoDB endpoint: {endpoint}"); + loader = loader.endpoint_url(endpoint); + } + if let (Some(access_key_id), Some(secret_access_key)) = + (&self.config.access_key_id, &self.config.secret_access_key) + { + info!( + "Using explicit DynamoDB credentials for sink ID: {}", + self.id + ); + loader = loader.credentials_provider(Credentials::new( + access_key_id.expose_secret(), + secret_access_key.expose_secret(), + self.config + .session_token + .as_ref() + .map(|token| token.expose_secret().to_owned()), + None, + CREDENTIALS_PROVIDER_NAME, + )); + } else { + info!( + "No explicit credentials provided, using the default AWS credential chain for sink ID: {}", + self.id + ); + } + + Ok(Client::new(&loader.load().await)) + } + + /// A key field that does not match the table makes DynamoDB reject every + /// write, and the runtime stops the connector on the first error, so the + /// mismatch is reported while the sink is still opening. + fn validate_key_schema(&self, key_schema: Vec) -> Result<(), Error> { + let table_key = |key_type: KeyType| { + key_schema + .iter() + .find(|element| element.key_type == key_type) + .map(|element| element.attribute_name.clone()) + }; + + let partition_key = table_key(KeyType::Hash); + if let Some(partition_key) = &partition_key + && partition_key != &self.partition_key_field + { + return Err(Error::InvalidConfigValue(format!( + "Table '{}' uses '{partition_key}' as its partition key, but partition_key_field is '{}'", + self.config.table, self.partition_key_field + ))); + } + + match (table_key(KeyType::Range), &self.sort_key_field) { + (Some(sort_key), Some(sort_key_field)) if &sort_key != sort_key_field => { + Err(Error::InvalidConfigValue(format!( + "Table '{}' uses '{sort_key}' as its sort key, but sort_key_field is '{sort_key_field}'", + self.config.table + ))) + } + (Some(sort_key), None) => Err(Error::InvalidConfigValue(format!( + "Table '{}' has a sort key '{sort_key}', so sort_key_field must be set", + self.config.table + ))), + (None, Some(sort_key_field)) => Err(Error::InvalidConfigValue(format!( + "Table '{}' has no sort key, so sort_key_field '{sort_key_field}' must be removed", + self.config.table + ))), + _ => Ok(()), + } + } + + async fn write_messages( + &self, + topic_metadata: &TopicMetadata, + messages_metadata: &MessagesMetadata, + mut messages: Vec, + ) -> Result<(), Error> { + let client = self.get_client()?; + let mut items = Vec::with_capacity(messages.len()); + let mut skipped = 0u64; + + for message in messages.iter_mut() { + match self.build_item(topic_metadata, messages_metadata, message) { + Ok(item) => items.push(item), + Err(reason) => { + skipped += 1; + warn!( + "DynamoDB sink ID: {} skipped message at offset: {}, reason: {reason}", + self.id, message.offset + ); + } + } + } + + let (mut items, duplicates) = self.deduplicate_items(items); + if duplicates > 0 { + debug!( + "DynamoDB sink ID: {} dropped {duplicates} items sharing a primary key", + self.id + ); + self.items_deduplicated + .fetch_add(duplicates, Ordering::Relaxed); + } + if skipped > 0 { + self.items_skipped.fetch_add(skipped, Ordering::Relaxed); + } + + let mut written = 0u64; + let mut last_error: Option = None; + while !items.is_empty() { + let chunk_size = self.batch_size.min(items.len()); + let chunk = items.drain(..chunk_size).collect::>(); + match self.write_chunk(client, chunk).await { + Ok(count) => written += count, + Err(error) => { + self.write_errors + .fetch_add(chunk_size as u64, Ordering::Relaxed); + error!( + "DynamoDB sink ID: {} failed to write {chunk_size} items to table: {}, error: {error}", + self.id, self.config.table + ); + last_error = Some(error); + } + } + } + self.items_written.fetch_add(written, Ordering::Relaxed); + + let table = &self.config.table; + if self.verbose { + info!( + "DynamoDB sink ID: {} wrote {written} items to table: {table}, current_offset: {}", + self.id, messages_metadata.current_offset + ); + } else { + debug!( + "DynamoDB sink ID: {} wrote {written} items to table: {table}, current_offset: {}", + self.id, messages_metadata.current_offset + ); + } + + match last_error { + Some(error) => Err(error), + None => Ok(()), + } + } + + /// `BatchWriteItem` rejects a request that carries two writes for the same + /// primary key, so only the newest item per key survives. DynamoDB would + /// overwrite the older ones anyway. + fn deduplicate_items( + &self, + items: Vec>, + ) -> (Vec>, u64) { + let mut positions: HashMap = HashMap::with_capacity(items.len()); + let mut deduplicated: Vec> = + Vec::with_capacity(items.len()); + let mut duplicates = 0u64; + + for item in items { + let key = self.key_signature(&item); + match positions.get(&key) { + Some(&position) => { + deduplicated[position] = item; + duplicates += 1; + } + None => { + positions.insert(key, deduplicated.len()); + deduplicated.push(item); + } + } + } + + (deduplicated, duplicates) + } + + /// `AttributeValue` is not hashable, so the key attributes are rendered + /// into a string that only lives for the deduplication pass. + fn key_signature(&self, item: &HashMap) -> String { + let mut signature = key_attribute_signature(item.get(&self.partition_key_field)); + if let Some(sort_key_field) = &self.sort_key_field { + signature.push('|'); + signature.push_str(&key_attribute_signature(item.get(sort_key_field))); + } + signature + } + + /// One `BatchWriteItem` round, retrying both throttled requests and the + /// `UnprocessedItems` the API returns inside an otherwise successful + /// response. + async fn write_chunk( + &self, + client: &Client, + items: Vec>, + ) -> Result { + let total = items.len() as u64; + let mut pending = Vec::with_capacity(items.len()); + for item in items { + let put_request = + PutRequest::builder() + .set_item(Some(item)) + .build() + .map_err(|error| { + Error::InvalidRecordValue(format!( + "Cannot build DynamoDB put request: {error}" + )) + })?; + pending.push(WriteRequest::builder().put_request(put_request).build()); + } + + let mut attempt = 0u32; + loop { + // `request_items` consumes the batch, and a failed send does not + // hand it back, so the retry loop keeps its own copy. + let result = client + .batch_write_item() + .request_items(&self.config.table, pending.clone()) + .send() + .await; + + let unprocessed = match result { + Ok(output) => output + .unprocessed_items + .and_then(|mut tables| tables.remove(&self.config.table)) + .unwrap_or_default(), + Err(error) => { + if !is_transient_error(&error) { + return Err(Error::PermanentHttpError(format!( + "DynamoDB batch write to table '{}' failed, error: {}", + self.config.table, + describe_sdk_error(&error) + ))); + } + attempt += 1; + if attempt > self.max_retries { + return Err(Error::CannotStoreData(format!( + "DynamoDB batch write to table '{}' failed after {attempt} attempts, error: {}", + self.config.table, + describe_sdk_error(&error) + ))); + } + warn!( + "Transient DynamoDB error on attempt {attempt}/{} for sink ID: {}, error: {}", + self.max_retries, + self.id, + describe_sdk_error(&error) + ); + self.backoff(attempt).await; + continue; + } + }; + + if unprocessed.is_empty() { + return Ok(total); + } + + attempt += 1; + if attempt > self.max_retries { + return Err(Error::CannotStoreData(format!( + "DynamoDB left {} of {total} items unprocessed in table '{}' after {attempt} attempts", + unprocessed.len(), + self.config.table + ))); + } + warn!( + "DynamoDB returned {} unprocessed items on attempt {attempt}/{} for sink ID: {}", + unprocessed.len(), + self.max_retries, + self.id + ); + self.backoff(attempt).await; + pending = unprocessed; + } + } + + async fn backoff(&self, attempt: u32) { + let delay = jitter(exponential_backoff( + self.retry_delay, + attempt - 1, + self.max_retry_delay, + )); + tokio::time::sleep(delay).await; + } + + fn build_item( + &self, + topic_metadata: &TopicMetadata, + messages_metadata: &MessagesMetadata, + message: &mut ConsumedMessage, + ) -> Result, Error> { + // The payload is moved out to build the item without copying it, so + // only the message metadata is read afterwards. + let payload = std::mem::replace(&mut message.payload, Payload::Raw(Vec::new())); + let mut item = payload_into_item(payload)?; + + if self.include_metadata { + item.insert( + "iggy_stream".to_owned(), + AttributeValue::S(topic_metadata.stream.clone()), + ); + item.insert( + "iggy_topic".to_owned(), + AttributeValue::S(topic_metadata.topic.clone()), + ); + item.insert( + "iggy_partition_id".to_owned(), + AttributeValue::N(messages_metadata.partition_id.to_string()), + ); + item.insert( + "iggy_offset".to_owned(), + AttributeValue::N(message.offset.to_string()), + ); + item.insert( + "iggy_timestamp".to_owned(), + AttributeValue::N(message.timestamp.to_string()), + ); + } + if self.include_checksum { + item.insert( + "iggy_checksum".to_owned(), + AttributeValue::N(message.checksum.to_string()), + ); + } + if self.include_origin_timestamp { + item.insert( + "iggy_origin_timestamp".to_owned(), + AttributeValue::N(message.origin_timestamp.to_string()), + ); + } + + item.entry(self.partition_key_field.clone()) + .or_insert_with(|| { + AttributeValue::S(build_message_key( + topic_metadata, + messages_metadata, + message.id, + )) + }); + if let Some(sort_key_field) = &self.sort_key_field { + item.entry(sort_key_field.clone()) + .or_insert_with(|| AttributeValue::N(message.offset.to_string())); + } + + validate_key_attribute(&item, &self.partition_key_field)?; + if let Some(sort_key_field) = &self.sort_key_field { + validate_key_attribute(&item, sort_key_field)?; + } + + let size = estimate_item_size(&item); + if size > self.max_item_size { + return Err(Error::InvalidRecordValue(format!( + "item size of {size} bytes exceeds the limit of {} bytes", + self.max_item_size + ))); + } + + Ok(item) + } + + fn get_client(&self) -> Result<&Client, Error> { + self.client + .as_ref() + .ok_or_else(|| Error::InitError("DynamoDB client is not connected".to_owned())) + } +} + +fn parse_duration(raw: Option<&str>, default: &str) -> Duration { + let value = raw.unwrap_or(default); + HumanDuration::from_str(value) + .map(Duration::from) + .unwrap_or_else(|_| { + warn!("Invalid DynamoDB sink duration '{value}', falling back to '{default}'"); + HumanDuration::from_str(default) + .map(Duration::from) + .unwrap_or(Duration::from_millis(500)) + }) +} + +/// DynamoDB items are attribute maps, so anything that is not a JSON object +/// is nested under a single payload attribute. +fn payload_into_item(payload: Payload) -> Result, Error> { + match payload { + Payload::Json(value) => Ok(match json_into_attribute_value(value) { + AttributeValue::M(item) => item, + other => HashMap::from([(PAYLOAD_FIELD.to_owned(), other)]), + }), + Payload::Text(text) => Ok(HashMap::from([( + PAYLOAD_FIELD.to_owned(), + AttributeValue::S(text), + )])), + Payload::Raw(bytes) => { + // simd-json unescapes in place, so a failed parse leaves the + // buffer rewritten and only the copy can be thrown away. + let mut buffer = bytes.clone(); + let attribute = match simd_json::to_owned_value(&mut buffer) { + Ok(value) => json_into_attribute_value(value), + Err(_) => AttributeValue::B(Blob::new(bytes)), + }; + Ok(match attribute { + AttributeValue::M(item) => item, + other => HashMap::from([(PAYLOAD_FIELD.to_owned(), other)]), + }) + } + Payload::Proto(_) | Payload::FlatBuffer(_) | Payload::Avro(_) => { + Err(Error::InvalidPayloadType) + } + } +} + +fn json_into_attribute_value(value: OwnedValue) -> AttributeValue { + match value { + OwnedValue::Static(StaticNode::Null) => AttributeValue::Null(true), + OwnedValue::Static(StaticNode::Bool(value)) => AttributeValue::Bool(value), + OwnedValue::Static(StaticNode::I64(value)) => AttributeValue::N(value.to_string()), + OwnedValue::Static(StaticNode::U64(value)) => AttributeValue::N(value.to_string()), + // DynamoDB numbers have no NaN or infinity, so those become null + // instead of a value the API would reject. + OwnedValue::Static(StaticNode::F64(value)) => { + if value.is_finite() { + AttributeValue::N(value.to_string()) + } else { + AttributeValue::Null(true) + } + } + OwnedValue::String(value) => AttributeValue::S(value), + OwnedValue::Array(values) => AttributeValue::L( + values + .into_iter() + .map(json_into_attribute_value) + .collect::>(), + ), + OwnedValue::Object(values) => AttributeValue::M( + values + .into_iter() + .map(|(key, value)| (key, json_into_attribute_value(value))) + .collect(), + ), + } +} + +fn key_attribute_signature(value: Option<&AttributeValue>) -> String { + match value { + Some(AttributeValue::S(text)) => format!("S:{text}"), + Some(AttributeValue::N(number)) => format!("N:{number}"), + Some(AttributeValue::B(blob)) => { + let mut signature = String::from("B:"); + for byte in blob.as_ref() { + let _ = write!(signature, "{byte:02x}"); + } + signature + } + _ => String::new(), + } +} + +fn build_message_key( + topic_metadata: &TopicMetadata, + messages_metadata: &MessagesMetadata, + message_id: u128, +) -> String { + format!( + "{}:{}:{}:{message_id}", + topic_metadata.stream, topic_metadata.topic, messages_metadata.partition_id + ) +} + +/// DynamoDB key attributes must be a non-empty string, number, or binary +/// value. Anything else fails the whole batch with a validation error, so the +/// offending message is dropped before it is sent. +fn validate_key_attribute( + item: &HashMap, + field: &str, +) -> Result<(), Error> { + let value = item + .get(field) + .ok_or_else(|| Error::InvalidRecordValue(format!("key field '{field}' is missing")))?; + let valid = match value { + AttributeValue::S(text) => !text.is_empty(), + AttributeValue::N(number) => !number.is_empty(), + AttributeValue::B(blob) => !blob.as_ref().is_empty(), + _ => false, + }; + if valid { + Ok(()) + } else { + Err(Error::InvalidRecordValue(format!( + "key field '{field}' must be a non-empty string, number, or binary value" + ))) + } +} + +fn estimate_item_size(item: &HashMap) -> usize { + item.iter() + .map(|(name, value)| name.len() + attribute_value_size(value)) + .sum() +} + +fn attribute_value_size(value: &AttributeValue) -> usize { + match value { + AttributeValue::S(text) => text.len(), + AttributeValue::N(number) => number.len(), + AttributeValue::B(blob) => blob.as_ref().len(), + AttributeValue::Bool(_) | AttributeValue::Null(_) => 1, + AttributeValue::Ss(values) => values.iter().map(String::len).sum(), + AttributeValue::Ns(values) => values.iter().map(String::len).sum(), + AttributeValue::Bs(values) => values.iter().map(|blob| blob.as_ref().len()).sum(), + AttributeValue::L(values) => values.iter().map(attribute_value_size).sum(), + AttributeValue::M(values) => estimate_item_size(values), + _ => 0, + } +} + +fn is_transient_error(error: &SdkError) -> bool +where + E: ProvideErrorMetadata, +{ + match error { + SdkError::TimeoutError(_) | SdkError::DispatchFailure(_) => true, + SdkError::ResponseError(_) => true, + SdkError::ServiceError(service_error) => { + is_transient_code(service_error.err().code().unwrap_or_default()) + } + _ => false, + } +} + +fn is_transient_code(code: &str) -> bool { + matches!( + code, + "ProvisionedThroughputExceededException" + | "RequestLimitExceeded" + | "ThrottlingException" + | "InternalServerError" + | "ServiceUnavailable" + | "TransactionInProgressException" + ) +} + +fn describe_sdk_error(error: &SdkError) -> String +where + E: ProvideErrorMetadata, +{ + match error.code() { + Some(code) => format!("{code}: {}", error.message().unwrap_or("no message")), + None => error.to_string(), + } +} + +#[cfg(test)] +mod tests { + use super::*; + use iggy_connector_sdk::Schema; + + fn given_default_config() -> DynamoDbSinkConfig { + DynamoDbSinkConfig { + table: "iggy_messages".to_owned(), + region: Some("us-east-1".to_owned()), + endpoint: None, + access_key_id: None, + secret_access_key: None, + session_token: None, + partition_key_field: None, + sort_key_field: None, + batch_size: None, + include_metadata: None, + include_checksum: None, + include_origin_timestamp: None, + max_item_size: None, + max_retries: None, + retry_delay: None, + max_retry_delay: None, + verbose_logging: None, + } + } + + fn given_topic_metadata() -> TopicMetadata { + TopicMetadata { + stream: "test_stream".to_owned(), + topic: "test_topic".to_owned(), + } + } + + fn given_messages_metadata() -> MessagesMetadata { + MessagesMetadata { + partition_id: 1, + current_offset: 10, + schema: Schema::Json, + } + } + + fn given_message(payload: Payload) -> ConsumedMessage { + ConsumedMessage { + id: 42, + offset: 7, + checksum: 99, + timestamp: 1_700_000_000, + origin_timestamp: 1_600_000_000, + headers: None, + payload, + } + } + + fn given_json_payload(raw: &str) -> Payload { + let mut bytes = raw.as_bytes().to_vec(); + Payload::Json(simd_json::to_owned_value(&mut bytes).expect("parse JSON")) + } + + #[test] + fn given_no_batch_size_when_created_should_use_dynamodb_limit() { + let sink = DynamoDbSink::new(1, given_default_config()); + assert_eq!(sink.batch_size, MAX_BATCH_WRITE_ITEMS); + } + + #[test] + fn given_oversized_batch_size_when_created_should_clamp_to_dynamodb_limit() { + let mut config = given_default_config(); + config.batch_size = Some(500); + let sink = DynamoDbSink::new(1, config); + assert_eq!(sink.batch_size, MAX_BATCH_WRITE_ITEMS); + } + + #[test] + fn given_zero_batch_size_when_created_should_use_single_item_batches() { + let mut config = given_default_config(); + config.batch_size = Some(0); + let sink = DynamoDbSink::new(1, config); + assert_eq!(sink.batch_size, 1); + } + + #[test] + fn given_reversed_retry_delays_when_created_should_raise_max_retry_delay() { + let mut config = given_default_config(); + config.retry_delay = Some("5s".to_owned()); + config.max_retry_delay = Some("1s".to_owned()); + let sink = DynamoDbSink::new(1, config); + assert_eq!(sink.retry_delay, Duration::from_secs(5)); + assert_eq!(sink.max_retry_delay, Duration::from_secs(5)); + } + + #[test] + fn given_invalid_duration_when_created_should_use_default() { + let mut config = given_default_config(); + config.retry_delay = Some("not-a-duration".to_owned()); + let sink = DynamoDbSink::new(1, config); + assert_eq!(sink.retry_delay, Duration::from_millis(500)); + } + + #[test] + fn given_oversized_max_item_size_when_created_should_clamp_to_dynamodb_limit() { + let mut config = given_default_config(); + config.max_item_size = Some(MAX_ITEM_SIZE * 2); + let sink = DynamoDbSink::new(1, config); + assert_eq!(sink.max_item_size, MAX_ITEM_SIZE); + } + + #[test] + fn given_json_object_payload_when_built_should_map_attributes() { + let sink = DynamoDbSink::new(1, given_default_config()); + let mut message = given_message(given_json_payload( + r#"{"name":"first","count":3,"ratio":1.5,"active":true,"tags":["a","b"],"nested":{"key":"value"},"missing":null}"#, + )); + + let item = sink + .build_item( + &given_topic_metadata(), + &given_messages_metadata(), + &mut message, + ) + .expect("build item"); + + assert_eq!(item["name"], AttributeValue::S("first".to_owned())); + assert_eq!(item["count"], AttributeValue::N("3".to_owned())); + assert_eq!(item["ratio"], AttributeValue::N("1.5".to_owned())); + assert_eq!(item["active"], AttributeValue::Bool(true)); + assert_eq!(item["missing"], AttributeValue::Null(true)); + assert_eq!( + item["tags"], + AttributeValue::L(vec![ + AttributeValue::S("a".to_owned()), + AttributeValue::S("b".to_owned()) + ]) + ); + assert_eq!( + item["nested"], + AttributeValue::M(HashMap::from([( + "key".to_owned(), + AttributeValue::S("value".to_owned()) + )])) + ); + } + + #[test] + fn given_json_object_payload_when_built_should_add_metadata_attributes() { + let sink = DynamoDbSink::new(1, given_default_config()); + let mut message = given_message(given_json_payload(r#"{"name":"first"}"#)); + + let item = sink + .build_item( + &given_topic_metadata(), + &given_messages_metadata(), + &mut message, + ) + .expect("build item"); + + assert_eq!( + item["iggy_stream"], + AttributeValue::S("test_stream".to_owned()) + ); + assert_eq!( + item["iggy_topic"], + AttributeValue::S("test_topic".to_owned()) + ); + assert_eq!(item["iggy_partition_id"], AttributeValue::N("1".to_owned())); + assert_eq!(item["iggy_offset"], AttributeValue::N("7".to_owned())); + assert_eq!(item["iggy_checksum"], AttributeValue::N("99".to_owned())); + assert_eq!( + item["iggy_origin_timestamp"], + AttributeValue::N("1600000000".to_owned()) + ); + assert_eq!( + item[DEFAULT_PARTITION_KEY_FIELD], + AttributeValue::S("test_stream:test_topic:1:42".to_owned()) + ); + } + + #[test] + fn given_disabled_metadata_when_built_should_only_keep_payload_and_key() { + let mut config = given_default_config(); + config.include_metadata = Some(false); + config.include_checksum = Some(false); + config.include_origin_timestamp = Some(false); + let sink = DynamoDbSink::new(1, config); + let mut message = given_message(given_json_payload(r#"{"name":"first"}"#)); + + let item = sink + .build_item( + &given_topic_metadata(), + &given_messages_metadata(), + &mut message, + ) + .expect("build item"); + + assert_eq!(item.len(), 2); + assert!(item.contains_key("name")); + assert!(item.contains_key(DEFAULT_PARTITION_KEY_FIELD)); + } + + #[test] + fn given_payload_with_partition_key_when_built_should_keep_payload_value() { + let mut config = given_default_config(); + config.partition_key_field = Some("user_id".to_owned()); + let sink = DynamoDbSink::new(1, config); + let mut message = given_message(given_json_payload(r#"{"user_id":"u-1"}"#)); + + let item = sink + .build_item( + &given_topic_metadata(), + &given_messages_metadata(), + &mut message, + ) + .expect("build item"); + + assert_eq!(item["user_id"], AttributeValue::S("u-1".to_owned())); + } + + #[test] + fn given_missing_sort_key_when_built_should_inject_offset() { + let mut config = given_default_config(); + config.sort_key_field = Some("event_offset".to_owned()); + let sink = DynamoDbSink::new(1, config); + let mut message = given_message(given_json_payload(r#"{"name":"first"}"#)); + + let item = sink + .build_item( + &given_topic_metadata(), + &given_messages_metadata(), + &mut message, + ) + .expect("build item"); + + assert_eq!(item["event_offset"], AttributeValue::N("7".to_owned())); + } + + #[test] + fn given_invalid_key_type_when_built_should_skip_message() { + let mut config = given_default_config(); + config.partition_key_field = Some("user_id".to_owned()); + let sink = DynamoDbSink::new(1, config); + let mut message = given_message(given_json_payload(r#"{"user_id":true}"#)); + + let result = sink.build_item( + &given_topic_metadata(), + &given_messages_metadata(), + &mut message, + ); + + assert!(result.is_err()); + } + + #[test] + fn given_oversized_payload_when_built_should_skip_message() { + let mut config = given_default_config(); + config.max_item_size = Some(64); + let sink = DynamoDbSink::new(1, config); + let mut message = given_message(Payload::Text("x".repeat(1024))); + + let result = sink.build_item( + &given_topic_metadata(), + &given_messages_metadata(), + &mut message, + ); + + assert!(result.is_err()); + } + + #[test] + fn given_unsupported_schema_payload_when_built_should_skip_message() { + let sink = DynamoDbSink::new(1, given_default_config()); + let mut message = given_message(Payload::Avro(vec![1, 2, 3])); + + let result = sink.build_item( + &given_topic_metadata(), + &given_messages_metadata(), + &mut message, + ); + + assert!(result.is_err()); + } + + #[test] + fn given_text_payload_when_built_should_store_payload_attribute() { + let sink = DynamoDbSink::new(1, given_default_config()); + let mut message = given_message(Payload::Text("hello".to_owned())); + + let item = sink + .build_item( + &given_topic_metadata(), + &given_messages_metadata(), + &mut message, + ) + .expect("build item"); + + assert_eq!(item[PAYLOAD_FIELD], AttributeValue::S("hello".to_owned())); + } + + #[test] + fn given_raw_json_payload_when_built_should_map_attributes() { + let sink = DynamoDbSink::new(1, given_default_config()); + let mut message = given_message(Payload::Raw(br#"{"name":"first"}"#.to_vec())); + + let item = sink + .build_item( + &given_topic_metadata(), + &given_messages_metadata(), + &mut message, + ) + .expect("build item"); + + assert_eq!(item["name"], AttributeValue::S("first".to_owned())); + } + + #[test] + fn given_raw_binary_payload_when_built_should_store_binary_attribute() { + let sink = DynamoDbSink::new(1, given_default_config()); + let mut message = given_message(Payload::Raw(vec![0xff, 0xfe, 0xfd])); + + let item = sink + .build_item( + &given_topic_metadata(), + &given_messages_metadata(), + &mut message, + ) + .expect("build item"); + + assert_eq!( + item[PAYLOAD_FIELD], + AttributeValue::B(Blob::new(vec![0xff, 0xfe, 0xfd])) + ); + } + + #[test] + fn given_json_array_payload_when_built_should_nest_under_payload_attribute() { + let sink = DynamoDbSink::new(1, given_default_config()); + let mut message = given_message(given_json_payload(r#"[1,2]"#)); + + let item = sink + .build_item( + &given_topic_metadata(), + &given_messages_metadata(), + &mut message, + ) + .expect("build item"); + + assert_eq!( + item[PAYLOAD_FIELD], + AttributeValue::L(vec![ + AttributeValue::N("1".to_owned()), + AttributeValue::N("2".to_owned()) + ]) + ); + } + + #[test] + fn given_non_finite_json_number_when_converted_should_become_null() { + let value = OwnedValue::Static(StaticNode::F64(f64::NAN)); + assert_eq!(json_into_attribute_value(value), AttributeValue::Null(true)); + } + + #[test] + fn given_same_message_when_built_twice_should_produce_the_same_key() { + let sink = DynamoDbSink::new(1, given_default_config()); + let topic_metadata = given_topic_metadata(); + let messages_metadata = given_messages_metadata(); + let mut first = given_message(given_json_payload(r#"{"name":"first"}"#)); + let mut second = given_message(given_json_payload(r#"{"name":"first"}"#)); + + let first_item = sink + .build_item(&topic_metadata, &messages_metadata, &mut first) + .expect("build item"); + let second_item = sink + .build_item(&topic_metadata, &messages_metadata, &mut second) + .expect("build item"); + + assert_eq!( + first_item[DEFAULT_PARTITION_KEY_FIELD], + second_item[DEFAULT_PARTITION_KEY_FIELD] + ); + } + + #[test] + fn given_different_topics_when_built_should_produce_different_keys() { + let messages_metadata = given_messages_metadata(); + let first_topic = given_topic_metadata(); + let second_topic = TopicMetadata { + stream: "test_stream".to_owned(), + topic: "other_topic".to_owned(), + }; + + assert_ne!( + build_message_key(&first_topic, &messages_metadata, 42), + build_message_key(&second_topic, &messages_metadata, 42) + ); + } + + #[test] + fn given_items_sharing_a_key_when_deduplicated_should_keep_the_newest() { + let mut config = given_default_config(); + config.partition_key_field = Some("user_id".to_owned()); + let sink = DynamoDbSink::new(1, config); + let items = vec![ + HashMap::from([ + ("user_id".to_owned(), AttributeValue::S("u-1".to_owned())), + ("name".to_owned(), AttributeValue::S("old".to_owned())), + ]), + HashMap::from([ + ("user_id".to_owned(), AttributeValue::S("u-2".to_owned())), + ("name".to_owned(), AttributeValue::S("other".to_owned())), + ]), + HashMap::from([ + ("user_id".to_owned(), AttributeValue::S("u-1".to_owned())), + ("name".to_owned(), AttributeValue::S("new".to_owned())), + ]), + ]; + + let (deduplicated, duplicates) = sink.deduplicate_items(items); + + assert_eq!(duplicates, 1); + assert_eq!(deduplicated.len(), 2); + assert_eq!(deduplicated[0]["name"], AttributeValue::S("new".to_owned())); + assert_eq!( + deduplicated[1]["name"], + AttributeValue::S("other".to_owned()) + ); + } + + #[test] + fn given_items_sharing_a_partition_key_when_sort_key_differs_should_keep_both() { + let mut config = given_default_config(); + config.partition_key_field = Some("user_id".to_owned()); + config.sort_key_field = Some("event_offset".to_owned()); + let sink = DynamoDbSink::new(1, config); + let items = vec![ + HashMap::from([ + ("user_id".to_owned(), AttributeValue::S("u-1".to_owned())), + ("event_offset".to_owned(), AttributeValue::N("1".to_owned())), + ]), + HashMap::from([ + ("user_id".to_owned(), AttributeValue::S("u-1".to_owned())), + ("event_offset".to_owned(), AttributeValue::N("2".to_owned())), + ]), + ]; + + let (deduplicated, duplicates) = sink.deduplicate_items(items); + + assert_eq!(duplicates, 0); + assert_eq!(deduplicated.len(), 2); + } + + #[tokio::test] + async fn given_no_client_when_messages_consumed_should_return_error() { + let sink = DynamoDbSink::new(1, given_default_config()); + let messages = vec![given_message(given_json_payload(r#"{"name":"first"}"#))]; + + let result = sink + .write_messages( + &given_topic_metadata(), + &given_messages_metadata(), + messages, + ) + .await; + + assert!(result.is_err()); + assert_eq!(sink.items_written.load(Ordering::Relaxed), 0); + } + + fn given_key_schema(partition_key: &str, sort_key: Option<&str>) -> Vec { + let mut key_schema = vec![ + KeySchemaElement::builder() + .attribute_name(partition_key) + .key_type(KeyType::Hash) + .build() + .expect("build key schema"), + ]; + if let Some(sort_key) = sort_key { + key_schema.push( + KeySchemaElement::builder() + .attribute_name(sort_key) + .key_type(KeyType::Range) + .build() + .expect("build key schema"), + ); + } + key_schema + } + + #[test] + fn given_matching_key_schema_when_validated_should_accept_the_table() { + let mut config = given_default_config(); + config.sort_key_field = Some("iggy_offset".to_owned()); + let sink = DynamoDbSink::new(1, config); + + let result = sink.validate_key_schema(given_key_schema( + DEFAULT_PARTITION_KEY_FIELD, + Some("iggy_offset"), + )); + + assert!(result.is_ok()); + } + + #[test] + fn given_another_partition_key_when_validated_should_reject_the_table() { + let sink = DynamoDbSink::new(1, given_default_config()); + + let result = sink.validate_key_schema(given_key_schema("user_id", None)); + + assert!(matches!(result, Err(Error::InvalidConfigValue(_)))); + } + + #[test] + fn given_table_sort_key_without_configured_field_when_validated_should_reject_the_table() { + let sink = DynamoDbSink::new(1, given_default_config()); + + let result = sink.validate_key_schema(given_key_schema( + DEFAULT_PARTITION_KEY_FIELD, + Some("iggy_offset"), + )); + + assert!(matches!(result, Err(Error::InvalidConfigValue(_)))); + } + + #[test] + fn given_configured_sort_key_without_table_sort_key_when_validated_should_reject_the_table() { + let mut config = given_default_config(); + config.sort_key_field = Some("iggy_offset".to_owned()); + let sink = DynamoDbSink::new(1, config); + + let result = sink.validate_key_schema(given_key_schema(DEFAULT_PARTITION_KEY_FIELD, None)); + + assert!(matches!(result, Err(Error::InvalidConfigValue(_)))); + } + + #[test] + fn given_raw_payload_that_is_not_json_when_built_should_store_the_original_bytes() { + let sink = DynamoDbSink::new(1, given_default_config()); + let raw = br#"{"a":"x\ny"} trailing"#.to_vec(); + let mut message = given_message(Payload::Raw(raw.clone())); + + let item = sink + .build_item( + &given_topic_metadata(), + &given_messages_metadata(), + &mut message, + ) + .expect("build item"); + + assert_eq!(item[PAYLOAD_FIELD], AttributeValue::B(Blob::new(raw))); + } + + #[test] + fn given_service_error_codes_when_checked_should_only_retry_the_transient_ones() { + assert!(is_transient_code("ThrottlingException")); + assert!(is_transient_code("ProvisionedThroughputExceededException")); + assert!(is_transient_code("InternalServerError")); + assert!(!is_transient_code("ValidationException")); + assert!(!is_transient_code("ResourceNotFoundException")); + assert!(!is_transient_code("")); + } + + #[test] + fn given_binary_keys_when_signed_should_not_collide() { + let sink = DynamoDbSink::new(1, given_default_config()); + let first = HashMap::from([( + DEFAULT_PARTITION_KEY_FIELD.to_owned(), + AttributeValue::B(Blob::new(vec![0x01, 0x02])), + )]); + let second = HashMap::from([( + DEFAULT_PARTITION_KEY_FIELD.to_owned(), + AttributeValue::B(Blob::new(vec![0x01, 0x03])), + )]); + + assert_ne!(sink.key_signature(&first), sink.key_signature(&second)); + } +} diff --git a/core/integration/Cargo.toml b/core/integration/Cargo.toml index 2a4acaa090..c937df5e10 100644 --- a/core/integration/Cargo.toml +++ b/core/integration/Cargo.toml @@ -36,6 +36,8 @@ login-session = ["dep:zbus-secret-service-keyring-store"] arrow = { workspace = true } assert_cmd = { workspace = true } async-trait = { workspace = true } +aws-config = { workspace = true } +aws-sdk-dynamodb = { workspace = true } base64 = { workspace = true } bon = { workspace = true } bytemuck = { workspace = true } diff --git a/core/integration/tests/connectors/dynamodb/dynamodb_sink.rs b/core/integration/tests/connectors/dynamodb/dynamodb_sink.rs new file mode 100644 index 0000000000..5d645e8a04 --- /dev/null +++ b/core/integration/tests/connectors/dynamodb/dynamodb_sink.rs @@ -0,0 +1,216 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +use crate::connectors::fixtures::{DynamoDbOps, DynamoDbSinkFixture, DynamoDbSinkSortKeyFixture}; +use aws_sdk_dynamodb::types::AttributeValue; +use bytes::Bytes; +use iggy::prelude::{IggyMessage, Partitioning}; +use iggy_common::{Identifier, MessageClient}; +use integration::harness::{TestHarness, seeds}; +use integration::iggy_harness; +use std::collections::HashMap; +use std::time::Duration; +use tokio::time::sleep; + +#[iggy_harness( + server(connectors_runtime(config_path = "tests/connectors/dynamodb/sink.toml")), + seed = seeds::connector_stream +)] +async fn given_json_messages_when_sink_consumes_should_write_items( + harness: &TestHarness, + fixture: DynamoDbSinkFixture, +) { + let payloads = [ + serde_json::json!({"name": "first", "count": 1}), + serde_json::json!({"name": "second", "count": 2}), + ]; + send_messages(harness, &payloads).await; + + let items = fixture + .wait_for_items(payloads.len()) + .await + .expect("wait for DynamoDB items"); + + assert_eq!(items.len(), payloads.len()); + assert_eq!( + string_attribute(&items[0], "iggy_stream"), + seeds::names::STREAM + ); + assert_eq!( + string_attribute(&items[0], "iggy_topic"), + seeds::names::TOPIC + ); + assert!( + items + .iter() + .any(|item| string_attribute(item, "name") == "first") + ); + assert!( + items + .iter() + .any(|item| string_attribute(item, "name") == "second") + ); + assert!( + items + .iter() + .all(|item| !string_attribute(item, "iggy_id").is_empty()) + ); + assert!( + items + .iter() + .any(|item| number_attribute(item, "count") == "1") + ); +} + +#[iggy_harness( + server(connectors_runtime(config_path = "tests/connectors/dynamodb/sink.toml")), + seed = seeds::connector_stream +)] +async fn given_table_with_sort_key_when_sink_consumes_should_write_offset_as_sort_key( + harness: &TestHarness, + fixture: DynamoDbSinkSortKeyFixture, +) { + let payloads = [ + serde_json::json!({"name": "first"}), + serde_json::json!({"name": "second"}), + ]; + send_messages(harness, &payloads).await; + + let items = fixture + .wait_for_items(payloads.len()) + .await + .expect("wait for DynamoDB items"); + + let mut offsets = items + .iter() + .map(|item| number_attribute(item, "iggy_offset").to_owned()) + .collect::>(); + offsets.sort(); + assert_eq!(offsets, vec!["0".to_owned(), "1".to_owned()]); +} + +#[iggy_harness( + server(connectors_runtime(config_path = "tests/connectors/dynamodb/sink.toml")), + seed = seeds::connector_stream +)] +async fn given_more_messages_than_a_batch_write_when_sink_consumes_should_write_every_item( + harness: &TestHarness, + fixture: DynamoDbSinkFixture, +) { + let payloads = (0..30) + .map(|index| serde_json::json!({"name": format!("message-{index}"), "count": index})) + .collect::>(); + send_messages(harness, &payloads).await; + + let items = fixture + .wait_for_items(payloads.len()) + .await + .expect("wait for DynamoDB items"); + + assert_eq!(items.len(), payloads.len()); +} + +#[iggy_harness( + server(connectors_runtime(config_path = "tests/connectors/dynamodb/sink.toml")), + seed = seeds::connector_stream +)] +async fn given_the_same_message_ids_when_sink_consumes_twice_should_overwrite_the_items( + harness: &TestHarness, + fixture: DynamoDbSinkFixture, +) { + let first = [ + serde_json::json!({"name": "first"}), + serde_json::json!({"name": "second"}), + ]; + send_messages(harness, &first).await; + fixture + .wait_for_items(first.len()) + .await + .expect("wait for DynamoDB items"); + + let second = [ + serde_json::json!({"name": "third"}), + serde_json::json!({"name": "fourth"}), + ]; + send_messages(harness, &second).await; + + let items = wait_for_names(&fixture, &["third", "fourth"]).await; + assert_eq!(items.len(), first.len()); +} + +/// The second batch reuses the message ids of the first one, so the item count +/// stays the same and only the payload tells the two rounds apart. +async fn wait_for_names( + fixture: &DynamoDbSinkFixture, + names: &[&str], +) -> Vec> { + for _ in 0..100 { + let items = fixture.scan_items().await.expect("scan DynamoDB items"); + if names.iter().all(|name| { + items + .iter() + .any(|item| string_attribute(item, "name") == *name) + }) { + return items; + } + sleep(Duration::from_millis(100)).await; + } + + panic!("DynamoDB items never carried the names: {names:?}"); +} + +async fn send_messages(harness: &TestHarness, payloads: &[serde_json::Value]) { + let client = harness.root_client().await.unwrap(); + let stream_id: Identifier = seeds::names::STREAM.try_into().unwrap(); + let topic_id: Identifier = seeds::names::TOPIC.try_into().unwrap(); + + let mut messages = payloads + .iter() + .enumerate() + .map(|(i, payload)| { + IggyMessage::builder() + .id((i + 1) as u128) + .payload(Bytes::from(serde_json::to_vec(payload).expect("serialize"))) + .build() + .expect("build message") + }) + .collect::>(); + + client + .send_messages( + &stream_id, + &topic_id, + &Partitioning::partition_id(0), + &mut messages, + ) + .await + .expect("send messages"); +} + +fn string_attribute<'a>(item: &'a HashMap, field: &str) -> &'a str { + match item.get(field) { + Some(AttributeValue::S(value)) => value, + _ => "", + } +} + +fn number_attribute<'a>(item: &'a HashMap, field: &str) -> &'a str { + match item.get(field) { + Some(AttributeValue::N(value)) => value, + _ => "", + } +} diff --git a/core/integration/tests/connectors/dynamodb/mod.rs b/core/integration/tests/connectors/dynamodb/mod.rs new file mode 100644 index 0000000000..f63001ed32 --- /dev/null +++ b/core/integration/tests/connectors/dynamodb/mod.rs @@ -0,0 +1,18 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +mod dynamodb_sink; diff --git a/core/integration/tests/connectors/dynamodb/sink.toml b/core/integration/tests/connectors/dynamodb/sink.toml new file mode 100644 index 0000000000..8bab26cd28 --- /dev/null +++ b/core/integration/tests/connectors/dynamodb/sink.toml @@ -0,0 +1,20 @@ +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you under the Apache License, Version 2.0 (the +# "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, +# software distributed under the License is distributed on an +# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +# KIND, either express or implied. See the License for the +# specific language governing permissions and limitations +# under the License. + +[connectors] +config_type = "local" +config_dir = "../connectors/sinks/dynamodb_sink" diff --git a/core/integration/tests/connectors/fixtures/dynamodb/fixture.rs b/core/integration/tests/connectors/fixtures/dynamodb/fixture.rs new file mode 100644 index 0000000000..1cd6b0a7d3 --- /dev/null +++ b/core/integration/tests/connectors/fixtures/dynamodb/fixture.rs @@ -0,0 +1,318 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +use crate::connectors::fixtures; +use async_trait::async_trait; +use aws_config::BehaviorVersion; +use aws_sdk_dynamodb::Client; +use aws_sdk_dynamodb::config::{Credentials, Region}; +use aws_sdk_dynamodb::types::{ + AttributeDefinition, AttributeValue, BillingMode, KeySchemaElement, KeyType, + ScalarAttributeType, +}; +use integration::harness::{TestBinaryError, TestFixture, seeds}; +use std::collections::HashMap; +use std::time::Duration; +use testcontainers_modules::testcontainers::core::{IntoContainerPort, WaitFor}; +use testcontainers_modules::testcontainers::runners::AsyncRunner; +use testcontainers_modules::testcontainers::{ContainerAsync, GenericImage, ImageExt}; +use tokio::time::sleep; +use tracing::info; + +const DYNAMODB_IMAGE: &str = "amazon/dynamodb-local"; +const DYNAMODB_TAG: &str = "2.5.2"; +const DYNAMODB_PORT: u16 = 8000; +const DYNAMODB_READY_LOG: &str = "Initializing DynamoDB Local"; +const DYNAMODB_REGION: &str = "us-east-1"; +const DYNAMODB_ACCESS_KEY: &str = "test"; +const DYNAMODB_SECRET_KEY: &str = "test"; +const CREDENTIALS_PROVIDER_NAME: &str = "iggy-dynamodb-test"; +pub const TEST_TABLE: &str = "iggy_messages"; +const PARTITION_KEY: &str = "iggy_id"; +const SORT_KEY: &str = "iggy_offset"; +/// Bounds the wait for the DynamoDB Local API to answer after its startup log. +const TABLE_CREATE_ATTEMPTS: usize = 30; +const TABLE_CREATE_RETRY_DELAY: Duration = Duration::from_secs(1); +const POLL_ATTEMPTS: usize = 100; +const POLL_INTERVAL: Duration = Duration::from_millis(100); + +const ENV_SINK_PATH: &str = "IGGY_CONNECTORS_SINK_DYNAMODB_PATH"; +const ENV_SINK_STREAMS_0_STREAM: &str = "IGGY_CONNECTORS_SINK_DYNAMODB_STREAMS_0_STREAM"; +const ENV_SINK_STREAMS_0_TOPICS: &str = "IGGY_CONNECTORS_SINK_DYNAMODB_STREAMS_0_TOPICS"; +const ENV_SINK_STREAMS_0_SCHEMA: &str = "IGGY_CONNECTORS_SINK_DYNAMODB_STREAMS_0_SCHEMA"; +const ENV_SINK_TABLE: &str = "IGGY_CONNECTORS_SINK_DYNAMODB_PLUGIN_CONFIG_TABLE"; +const ENV_SINK_REGION: &str = "IGGY_CONNECTORS_SINK_DYNAMODB_PLUGIN_CONFIG_REGION"; +const ENV_SINK_ENDPOINT: &str = "IGGY_CONNECTORS_SINK_DYNAMODB_PLUGIN_CONFIG_ENDPOINT"; +const ENV_SINK_ACCESS_KEY_ID: &str = "IGGY_CONNECTORS_SINK_DYNAMODB_PLUGIN_CONFIG_ACCESS_KEY_ID"; +const ENV_SINK_SECRET_ACCESS_KEY: &str = + "IGGY_CONNECTORS_SINK_DYNAMODB_PLUGIN_CONFIG_SECRET_ACCESS_KEY"; +const ENV_SINK_SORT_KEY_FIELD: &str = "IGGY_CONNECTORS_SINK_DYNAMODB_PLUGIN_CONFIG_SORT_KEY_FIELD"; + +pub trait DynamoDbOps: Sync { + fn client(&self) -> &Client; + + fn scan_items( + &self, + ) -> impl std::future::Future< + Output = Result>, TestBinaryError>, + > + Send { + async move { + let output = self + .client() + .scan() + .table_name(TEST_TABLE) + .send() + .await + .map_err(|error| TestBinaryError::InvalidState { + message: format!("Failed to scan DynamoDB table: {error}"), + })?; + Ok(output.items.unwrap_or_default()) + } + } + + fn wait_for_items( + &self, + expected_count: usize, + ) -> impl std::future::Future< + Output = Result>, TestBinaryError>, + > + Send { + async move { + let mut last_count = 0usize; + for _ in 0..POLL_ATTEMPTS { + if let Ok(items) = self.scan_items().await { + last_count = items.len(); + if last_count >= expected_count { + return Ok(items); + } + } + sleep(POLL_INTERVAL).await; + } + + Err(TestBinaryError::InvalidState { + message: format!( + "Expected {expected_count} DynamoDB items, found {last_count} after {POLL_ATTEMPTS} attempts" + ), + }) + } + } +} + +pub struct DynamoDbSinkFixture { + #[allow(dead_code)] + container: ContainerAsync, + client: Client, + endpoint: String, +} + +impl DynamoDbOps for DynamoDbSinkFixture { + fn client(&self) -> &Client { + &self.client + } +} + +impl DynamoDbSinkFixture { + async fn start(fixture_type: &str, with_sort_key: bool) -> Result { + let container = GenericImage::new(DYNAMODB_IMAGE, DYNAMODB_TAG) + .with_exposed_port(DYNAMODB_PORT.tcp()) + .with_wait_for(WaitFor::message_on_stdout(DYNAMODB_READY_LOG)) + .with_container_name(fixtures::unique_container_name("dynamodb")) + .with_mapped_port(0, DYNAMODB_PORT.tcp()) + .start() + .await + .map_err(|error| TestBinaryError::FixtureSetup { + fixture_type: fixture_type.to_string(), + message: format!("Failed to start DynamoDB Local container: {error}"), + })?; + + let mapped_port = container + .ports() + .await + .map_err(|error| TestBinaryError::FixtureSetup { + fixture_type: fixture_type.to_string(), + message: format!("Failed to get ports: {error}"), + })? + .map_to_host_port_ipv4(DYNAMODB_PORT) + .ok_or_else(|| TestBinaryError::FixtureSetup { + fixture_type: fixture_type.to_string(), + message: "No mapping for DynamoDB port".to_string(), + })?; + + let endpoint = format!("http://localhost:{mapped_port}"); + info!("DynamoDB Local container available at {endpoint}"); + + let config = aws_config::defaults(BehaviorVersion::latest()) + .region(Region::new(DYNAMODB_REGION)) + .endpoint_url(&endpoint) + .credentials_provider(Credentials::new( + DYNAMODB_ACCESS_KEY, + DYNAMODB_SECRET_KEY, + None, + None, + CREDENTIALS_PROVIDER_NAME, + )) + .load() + .await; + let client = Client::new(&config); + + create_table(&client, fixture_type, with_sort_key).await?; + + Ok(Self { + container, + client, + endpoint, + }) + } + + fn base_envs(&self) -> HashMap { + HashMap::from([ + ( + ENV_SINK_PATH.to_string(), + "../../target/debug/libiggy_connector_dynamodb_sink".to_string(), + ), + ( + ENV_SINK_STREAMS_0_STREAM.to_string(), + seeds::names::STREAM.to_string(), + ), + ( + ENV_SINK_STREAMS_0_TOPICS.to_string(), + format!("[{}]", seeds::names::TOPIC), + ), + (ENV_SINK_STREAMS_0_SCHEMA.to_string(), "json".to_string()), + (ENV_SINK_TABLE.to_string(), TEST_TABLE.to_string()), + (ENV_SINK_REGION.to_string(), DYNAMODB_REGION.to_string()), + (ENV_SINK_ENDPOINT.to_string(), self.endpoint.clone()), + ( + ENV_SINK_ACCESS_KEY_ID.to_string(), + DYNAMODB_ACCESS_KEY.to_string(), + ), + ( + ENV_SINK_SECRET_ACCESS_KEY.to_string(), + DYNAMODB_SECRET_KEY.to_string(), + ), + ]) + } +} + +#[async_trait] +impl TestFixture for DynamoDbSinkFixture { + async fn setup() -> Result { + Self::start("DynamoDbSinkFixture", false).await + } + + fn connectors_runtime_envs(&self) -> HashMap { + self.base_envs() + } +} + +/// Same container, but the table carries a sort key so the sink has to fill it +/// from the message offset. +pub struct DynamoDbSinkSortKeyFixture { + inner: DynamoDbSinkFixture, +} + +impl DynamoDbOps for DynamoDbSinkSortKeyFixture { + fn client(&self) -> &Client { + self.inner.client() + } +} + +#[async_trait] +impl TestFixture for DynamoDbSinkSortKeyFixture { + async fn setup() -> Result { + let inner = DynamoDbSinkFixture::start("DynamoDbSinkSortKeyFixture", true).await?; + Ok(Self { inner }) + } + + fn connectors_runtime_envs(&self) -> HashMap { + let mut envs = self.inner.base_envs(); + envs.insert(ENV_SINK_SORT_KEY_FIELD.to_string(), SORT_KEY.to_string()); + envs + } +} + +async fn create_table( + client: &Client, + fixture_type: &str, + with_sort_key: bool, +) -> Result<(), TestBinaryError> { + let partition_key_definition = AttributeDefinition::builder() + .attribute_name(PARTITION_KEY) + .attribute_type(ScalarAttributeType::S) + .build() + .map_err(|error| TestBinaryError::FixtureSetup { + fixture_type: fixture_type.to_string(), + message: format!("Failed to build attribute definition: {error}"), + })?; + let partition_key_schema = KeySchemaElement::builder() + .attribute_name(PARTITION_KEY) + .key_type(KeyType::Hash) + .build() + .map_err(|error| TestBinaryError::FixtureSetup { + fixture_type: fixture_type.to_string(), + message: format!("Failed to build key schema: {error}"), + })?; + + let mut last_error = String::new(); + for _ in 0..TABLE_CREATE_ATTEMPTS { + let mut request = client + .create_table() + .table_name(TEST_TABLE) + .billing_mode(BillingMode::PayPerRequest) + .attribute_definitions(partition_key_definition.clone()) + .key_schema(partition_key_schema.clone()); + + if with_sort_key { + let sort_key_definition = AttributeDefinition::builder() + .attribute_name(SORT_KEY) + .attribute_type(ScalarAttributeType::N) + .build() + .map_err(|error| TestBinaryError::FixtureSetup { + fixture_type: fixture_type.to_string(), + message: format!("Failed to build attribute definition: {error}"), + })?; + let sort_key_schema = KeySchemaElement::builder() + .attribute_name(SORT_KEY) + .key_type(KeyType::Range) + .build() + .map_err(|error| TestBinaryError::FixtureSetup { + fixture_type: fixture_type.to_string(), + message: format!("Failed to build key schema: {error}"), + })?; + request = request + .attribute_definitions(sort_key_definition) + .key_schema(sort_key_schema); + } + + match request.send().await { + Ok(_) => { + info!("DynamoDB table '{TEST_TABLE}' created"); + return Ok(()); + } + Err(error) => { + last_error = error.to_string(); + sleep(TABLE_CREATE_RETRY_DELAY).await; + } + } + } + + Err(TestBinaryError::FixtureSetup { + fixture_type: fixture_type.to_string(), + message: format!( + "Table '{TEST_TABLE}' not creatable after {TABLE_CREATE_ATTEMPTS} attempts, last error: {last_error}" + ), + }) +} diff --git a/core/integration/tests/connectors/fixtures/dynamodb/mod.rs b/core/integration/tests/connectors/fixtures/dynamodb/mod.rs new file mode 100644 index 0000000000..24fc43fdf5 --- /dev/null +++ b/core/integration/tests/connectors/fixtures/dynamodb/mod.rs @@ -0,0 +1,20 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +mod fixture; + +pub use fixture::{DynamoDbOps, DynamoDbSinkFixture, DynamoDbSinkSortKeyFixture}; diff --git a/core/integration/tests/connectors/fixtures/mod.rs b/core/integration/tests/connectors/fixtures/mod.rs index 7eaf6fe510..a8e0a200a6 100644 --- a/core/integration/tests/connectors/fixtures/mod.rs +++ b/core/integration/tests/connectors/fixtures/mod.rs @@ -20,6 +20,7 @@ use uuid::Uuid; mod clickhouse; mod delta; mod doris; +mod dynamodb; mod elasticsearch; mod http; mod iceberg; @@ -57,6 +58,7 @@ pub use doris::{ DorisOps, DorisSinkColumnsMappingFixture, DorisSinkCsvFixture, DorisSinkFixture, DorisSinkMaxFilterRatioFixture, DorisSinkPreCreatedFixture, }; +pub use dynamodb::{DynamoDbOps, DynamoDbSinkFixture, DynamoDbSinkSortKeyFixture}; pub use elasticsearch::{ElasticsearchSinkFixture, ElasticsearchSourcePreCreatedFixture}; pub use http::{ HttpSinkIndividualFixture, HttpSinkJsonArrayFixture, HttpSinkMultiTopicFixture, diff --git a/core/integration/tests/connectors/mod.rs b/core/integration/tests/connectors/mod.rs index f6346bcc20..52f96b7c8f 100644 --- a/core/integration/tests/connectors/mod.rs +++ b/core/integration/tests/connectors/mod.rs @@ -19,6 +19,7 @@ mod api; mod clickhouse; mod delta; mod doris; +mod dynamodb; mod elasticsearch; mod fixtures; mod http;