This is an automated email from the ASF dual-hosted git repository.
github-merge-queue[bot] pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/seatunnel.git
The following commit(s) were added to refs/heads/dev by this push:
new 96048a7b81 [Feature][Connector-V2] Add Azure CosmosDB Source-Side
Connector (#11167)
96048a7b81 is described below
commit 96048a7b81c92d743acc23b5bed652bdfef9d6e5
Author: Tony Nguyen <[email protected]>
AuthorDate: Sun Aug 30 16:30:36 2026 +0000
[Feature][Connector-V2] Add Azure CosmosDB Source-Side Connector (#11167)
---
.github/workflows/labeler/label-scope-conf.yml | 5 +
config/plugin_config | 1 +
.../changelog/connector-azurecosmosdb.md | 7 +
docs/en/connectors/source/AzureCosmosDB.md | 237 ++++++++++++++++++++
.../changelog/connector-azurecosmosdb.md | 7 +
docs/zh/connectors/source/AzureCosmosDB.md | 237 ++++++++++++++++++++
plugin-mapping.properties | 1 +
.../connector-azurecosmosdb/pom.xml | 45 ++++
.../azurecosmosdb/config/AzureCosmosDBConfig.java | 181 +++++++++++++++
.../config/AzureCosmosDBSourceOptions.java | 93 ++++++++
.../serialize/CosmosItemDeserializer.java | 154 +++++++++++++
.../azurecosmosdb/source/AzureCosmosDBSource.java | 81 +++++++
.../source/AzureCosmosDBSourceFactory.java | 86 +++++++
.../source/AzureCosmosDBSourceReader.java | 242 ++++++++++++++++++++
.../source/AzureCosmosDBSourceSplit.java | 58 +++++
.../source/AzureCosmosDBSourceSplitEnumerator.java | 145 ++++++++++++
.../source/AzureCosmosDBSourceState.java | 44 ++++
.../AzureCosmosDBSourceFactoryTest.java | 39 ++++
.../azurecosmosdb/CosmosItemDeserializerTest.java | 72 ++++++
.../config/AzureCosmosDBConfigTest.java | 69 ++++++
.../source/AzureCosmosDBSourceReaderTest.java | 210 ++++++++++++++++++
.../AzureCosmosDBSourceSplitEnumeratorTest.java | 121 ++++++++++
seatunnel-connectors-v2/pom.xml | 1 +
seatunnel-dist/pom.xml | 6 +
.../connector-azurecosmosdb-e2e/pom.xml | 63 ++++++
.../azurecosmosdb/AbstractAzureCosmosDBIT.java | 247 +++++++++++++++++++++
.../azurecosmosdb/AzureCosmosDBSourceIT.java | 138 ++++++++++++
.../azurecosmosdb/azurecosmosdb_source_basic.conf | 73 ++++++
.../azurecosmosdb_source_pagination.conf | 68 ++++++
.../azurecosmosdb_source_query_filter.conf | 68 ++++++
seatunnel-e2e/seatunnel-connector-v2-e2e/pom.xml | 1 +
31 files changed, 2800 insertions(+)
diff --git a/.github/workflows/labeler/label-scope-conf.yml
b/.github/workflows/labeler/label-scope-conf.yml
index f405ddbf5e..e4cbc53bd3 100644
--- a/.github/workflows/labeler/label-scope-conf.yml
+++ b/.github/workflows/labeler/label-scope-conf.yml
@@ -84,6 +84,11 @@ amazonsqs:
- changed-files:
- any-glob-to-any-file:
seatunnel-connectors-v2/connector-amazonsqs/**
- all-globs-to-all-files:
'!seatunnel-connectors-v2/connector-!(amazonsqs)/**'
+azurecosmosdb:
+ - all:
+ - changed-files:
+ - any-glob-to-any-file:
seatunnel-connectors-v2/connector-azurecosmosdb/**
+ - all-globs-to-all-files:
'!seatunnel-connectors-v2/connector-!(azurecosmosdb)/**'
cassandra:
- all:
- changed-files:
diff --git a/config/plugin_config b/config/plugin_config
index 55b27d764e..4e2ab94100 100644
--- a/config/plugin_config
+++ b/config/plugin_config
@@ -21,6 +21,7 @@
# Don't modify the delimiter " -- ", just select the plugin you need
--connectors-v2--
connector-amazondynamodb
+connector-azurecosmosdb
connector-assert
connector-cassandra
connector-salesforce
diff --git a/docs/en/connectors/changelog/connector-azurecosmosdb.md
b/docs/en/connectors/changelog/connector-azurecosmosdb.md
new file mode 100644
index 0000000000..ae42106c4b
--- /dev/null
+++ b/docs/en/connectors/changelog/connector-azurecosmosdb.md
@@ -0,0 +1,7 @@
+<details><summary> Change Log </summary>
+
+| Change | Commit | Version |
+| --- | --- | --- |
+|[Improve][Connector-V2][AzureCosmosDB] Support source checkpoint resume for
paginated reads|-|Next|
+
+</details>
diff --git a/docs/en/connectors/source/AzureCosmosDB.md
b/docs/en/connectors/source/AzureCosmosDB.md
new file mode 100644
index 0000000000..b21f5914c5
--- /dev/null
+++ b/docs/en/connectors/source/AzureCosmosDB.md
@@ -0,0 +1,237 @@
+import ChangeLog from '../changelog/connector-azurecosmosdb.md';
+
+# AzureCosmosDB
+
+> Azure Cosmos DB source connector
+
+## Support Connector Version
+
+- Azure Cosmos DB SQL (Core) API accounts
+
+## Support Those Engines
+
+> Spark<br/>
+> Flink<br/>
+> SeaTunnel Zeta<br/>
+
+## Key Features
+
+- [x] [batch](../../introduction/concepts/connector-v2-features.md)
+- [ ] [stream](../../introduction/concepts/connector-v2-features.md)
+- [ ] [exactly-once](../../introduction/concepts/connector-v2-features.md)
+- [ ] [column projection](../../introduction/concepts/connector-v2-features.md)
+- [ ] [parallelism](../../introduction/concepts/connector-v2-features.md)
+- [ ] [support user-defined
split](../../introduction/concepts/connector-v2-features.md)
+
+## Description
+
+Read data from Azure Cosmos DB (SQL API) containers using bounded batch scans.
+
+V1 scope for this connector:
+
+- source only
+- bounded/batch reads with a Cosmos SQL query
+- schema required (no schema inference)
+- catalog discovery out of scope
+- single split reading; physical partition/range parallel reads are out of
scope
+- change feed, managed identity, and container creation are out of scope
+
+## Supported DataSource Info
+
+In order to use the AzureCosmosDB connector, the following dependency is
required.
+It can be installed via `install-plugin.sh` or downloaded from the Maven
Central Repository.
+
+| Datasource | Supported Versions | Dependency
|
+|----------------|---------------------------------|----------------------------------------------------------------------------------------------------|
+| AzureCosmosDB | SQL (Core) API |
[Download](https://mvnrepository.com/artifact/org.apache.seatunnel/connector-azurecosmosdb)
|
+
+## Database Dependency
+
+> Please install the connector plugin before running jobs:
+
+```shell
+sh bin/install-plugin.sh ${version}
+```
+
+Ensure `connector-azurecosmosdb` is included in your plugin installation. The
connector uses the Azure Cosmos Java SDK (`azure-cosmos` 4.63.0) at runtime.
+
+## Data Type Mapping
+
+Cosmos DB stores JSON documents. You must define the SeaTunnel schema
explicitly. The connector maps JSON values to SeaTunnel types according to your
configured field types:
+
+| Cosmos JSON value | SeaTunnel Data type (when configured) |
+|-------------------|---------------------------------------|
+| Boolean | BOOLEAN |
+| Number | TINYINT / SMALLINT / INT / BIGINT / FLOAT / DOUBLE |
+| String | STRING |
+| String (ISO-8601) | DATE / TIME / TIMESTAMP |
+| String (decimal) | DECIMAL |
+| Binary | BYTES |
+| Object | MAP / ROW |
+| Array | ARRAY |
+| Null | null |
+
+## Source Options
+
+| Name | Type | Required | Default | Description |
+| --- | --- | --- | --- | --- |
+| uri | string | no | - | Azure Cosmos DB account URI |
+| endpoint | string | no | - | Azure Cosmos DB account endpoint |
+| key | string | no | - | Azure Cosmos DB account key |
+| primary_key | string | no | - | Azure Cosmos DB primary account key |
+| secondary_key | string | no | - | Azure Cosmos DB secondary account key |
+| primary_connection_string | string | no | - | Azure Cosmos DB primary
connection string |
+| secondary_connection_string | string | no | - | Azure Cosmos DB secondary
connection string |
+| database | string | yes | - | Azure Cosmos DB database name |
+| container | string | yes | - | Azure Cosmos DB container name |
+| schema | config | yes | - | Data schema definition |
+| query | string | no | SELECT * FROM c | Cosmos SQL query used to read source
data |
+| max_item_count | int | no | 100 | Max item count per query page |
+| common-options | - | no | - | Source plugin common parameters, see [Source
Common Options](../common-options/source-common-options.md) |
+
+### uri [string]
+
+Azure Cosmos DB account URI. Treated as an alias of `endpoint`.
+
+### endpoint [string]
+
+Azure Cosmos DB account endpoint, for example
`https://example-account.documents.azure.com:443/`.
+
+### key [string]
+
+Azure Cosmos DB account key.
+
+### primary_key [string]
+
+Azure Cosmos DB primary account key.
+
+### secondary_key [string]
+
+Azure Cosmos DB secondary account key.
+
+### primary_connection_string [string]
+
+Azure Cosmos DB primary connection string. The connector can parse
`AccountEndpoint` and `AccountKey` from it.
+
+### secondary_connection_string [string]
+
+Azure Cosmos DB secondary connection string. The connector can parse
`AccountEndpoint` and `AccountKey` from it.
+
+### database [string]
+
+Target database name. The database must already exist.
+
+### container [string]
+
+Target container name. The container must already exist.
+
+### schema [config]
+
+Cosmos DB stores JSON documents and does not enforce your SeaTunnel schema.
You must provide `schema.fields` explicitly. For details, refer to [Schema
Feature](../../introduction/concepts/schema-feature.md).
+
+Example:
+
+```hocon
+schema = {
+ fields {
+ id = string
+ user_id = string
+ amount = double
+ created_at = timestamp
+ labels = "map<string,string>"
+ }
+}
+```
+
+### query [string]
+
+Cosmos SQL query used for bounded batch reads. Use this for filtering and
field projection, for example `SELECT c.id, c.name FROM c WHERE c.score > 10`.
+
+### max_item_count [int]
+
+Preferred page size when the SDK iterates query results. During checkpoint
restore, the connector persists the SDK continuation token for the in-flight
paginated query and resumes from the last completed page.
+
+### common-options
+
+Source plugin common parameters, refer to [Source Common
Options](../common-options/source-common-options.md) for details.
+
+### Tips
+
+> 1. You must provide at least one of `uri`, `endpoint`,
`primary_connection_string`, or `secondary_connection_string`, and at least one
of `key`, `primary_key`, `secondary_key`, or a connection string that contains
an account key.<br/>
+> 2. V1 uses a single split. Increasing source parallelism does not
parallelize Cosmos reads across physical partitions.<br/>
+> 3. Use Cosmos SQL in `query` for filtering and projection. Connector-level
column projection is not supported as a separate feature.<br/>
+> 4. The connector reads from an existing container only. It does not create
databases, containers, or indexes.
+> 5. Checkpoint resume is based on Cosmos query page boundaries, not
individual rows. Change feed reading is still out of scope.
+
+## How to Create an Azure Cosmos DB Data Synchronization Job
+
+The following example demonstrates how to create a batch job that reads data
from Azure Cosmos DB and prints it to the local client:
+
+```bash
+# Set the basic configuration of the task to be performed
+env {
+ parallelism = 1
+ job.mode = "BATCH"
+}
+
+source {
+ AzureCosmosDB {
+ uri = "https://example-account.documents.azure.com:443/"
+ primary_key = "<cosmos-account-key>"
+ database = "app-db"
+ container = "orders"
+ query = "SELECT c.id, c.user_id, c.amount, c.created_at, c.labels FROM c"
+ max_item_count = 200
+ schema = {
+ fields {
+ id = string
+ user_id = string
+ amount = double
+ created_at = timestamp
+ labels = "map<string,string>"
+ }
+ }
+ }
+}
+
+sink {
+ Console {
+ parallelism = 1
+ }
+}
+```
+
+### Query with filter and smaller page size
+
+```bash
+env {
+ parallelism = 1
+ job.mode = "BATCH"
+}
+
+source {
+ AzureCosmosDB {
+ endpoint = "https://example-account.documents.azure.com:443/"
+ key = "<cosmos-account-key>"
+ database = "app-db"
+ container = "orders"
+ query = "SELECT c.id, c.user_id, c.amount FROM c WHERE c.amount > 100"
+ max_item_count = 50
+ schema = {
+ fields {
+ id = string
+ user_id = string
+ amount = double
+ }
+ }
+ }
+}
+
+sink {
+ Console {}
+}
+```
+
+## Changelog
+
+<ChangeLog />
diff --git a/docs/zh/connectors/changelog/connector-azurecosmosdb.md
b/docs/zh/connectors/changelog/connector-azurecosmosdb.md
new file mode 100644
index 0000000000..ae42106c4b
--- /dev/null
+++ b/docs/zh/connectors/changelog/connector-azurecosmosdb.md
@@ -0,0 +1,7 @@
+<details><summary> Change Log </summary>
+
+| Change | Commit | Version |
+| --- | --- | --- |
+|[Improve][Connector-V2][AzureCosmosDB] Support source checkpoint resume for
paginated reads|-|Next|
+
+</details>
diff --git a/docs/zh/connectors/source/AzureCosmosDB.md
b/docs/zh/connectors/source/AzureCosmosDB.md
new file mode 100644
index 0000000000..66ee102014
--- /dev/null
+++ b/docs/zh/connectors/source/AzureCosmosDB.md
@@ -0,0 +1,237 @@
+import ChangeLog from '../changelog/connector-azurecosmosdb.md';
+
+# AzureCosmosDB
+
+> Azure Cosmos DB 源连接器
+
+## 连接器支持版本
+
+- Azure Cosmos DB SQL(Core)API 账户
+
+## 支持这些引擎
+
+> Spark<br/>
+> Flink<br/>
+> SeaTunnel Zeta<br/>
+
+## 关键特性
+
+- [x] [批处理](../../introduction/concepts/connector-v2-features.md)
+- [ ] [流处理](../../introduction/concepts/connector-v2-features.md)
+- [ ] [精确一次](../../introduction/concepts/connector-v2-features.md)
+- [ ] [列投影](../../introduction/concepts/connector-v2-features.md)
+- [ ] [并行度](../../introduction/concepts/connector-v2-features.md)
+- [ ] [支持用户自定义分片](../../introduction/concepts/connector-v2-features.md)
+
+## 描述
+
+通过有界批处理扫描,从 Azure Cosmos DB(SQL API)容器读取数据。
+
+本连接器 V1 范围:
+
+- 仅 source
+- 使用 Cosmos SQL 的有界/批处理读取
+- 必须显式配置 schema(不支持 schema 推断)
+- 不支持 catalog 发现
+- 单分片读取;物理分区/范围并行读取暂不支持
+- change feed、托管身份认证、容器创建均不在 V1 范围内
+
+## 支持的数据源信息
+
+使用 AzureCosmosDB 连接器需要以下依赖。
+可通过 `install-plugin.sh` 安装,或从 Maven Central 下载。
+
+| 数据源 | 支持的版本 | 依赖
|
+|----------------|----------------------|----------------------------------------------------------------------------------------------|
+| AzureCosmosDB | SQL(Core)API |
[Download](https://mvnrepository.com/artifact/org.apache.seatunnel/connector-azurecosmosdb)
|
+
+## 数据库依赖
+
+> 运行作业前请先安装连接器插件:
+
+```shell
+sh bin/install-plugin.sh ${version}
+```
+
+请确保安装 `connector-azurecosmosdb` 插件。连接器运行时使用 Azure Cosmos Java
SDK(`azure-cosmos` 4.63.0)。
+
+## 数据类型映射
+
+Cosmos DB 存储 JSON 文档。你必须显式定义 SeaTunnel schema。连接器会根据你配置的字段类型,将 JSON 值映射为
SeaTunnel 类型:
+
+| Cosmos JSON 值 | SeaTunnel 数据类型(按 schema 配置) |
+|----------------|--------------------------------------|
+| Boolean | BOOLEAN |
+| Number | TINYINT / SMALLINT / INT / BIGINT / FLOAT / DOUBLE |
+| String | STRING |
+| String (ISO-8601) | DATE / TIME / TIMESTAMP |
+| String (decimal) | DECIMAL |
+| Binary | BYTES |
+| Object | MAP / ROW |
+| Array | ARRAY |
+| Null | null |
+
+## 源配置项
+
+| 名称 | 类型 | 必需 | 默认值 | 描述 |
+| --- | --- | --- | --- | --- |
+| uri | string | 否 | - | Azure Cosmos DB 账号 URI |
+| endpoint | string | 否 | - | Azure Cosmos DB 账号 endpoint |
+| key | string | 否 | - | Azure Cosmos DB 账号 key |
+| primary_key | string | 否 | - | Azure Cosmos DB 主账号 key |
+| secondary_key | string | 否 | - | Azure Cosmos DB 次账号 key |
+| primary_connection_string | string | 否 | - | Azure Cosmos DB 主连接字符串 |
+| secondary_connection_string | string | 否 | - | Azure Cosmos DB 次连接字符串 |
+| database | string | 是 | - | Azure Cosmos DB 数据库名 |
+| container | string | 是 | - | Azure Cosmos DB 容器名 |
+| schema | config | 是 | - | 数据 schema 定义 |
+| query | string | 否 | SELECT * FROM c | 读取数据使用的 Cosmos SQL |
+| max_item_count | int | 否 | 100 | 每页查询最大记录数 |
+| common-options | - | 否 | - | 源插件通用参数,详见 [Source Common
Options](../common-options/source-common-options.md) |
+
+### uri [string]
+
+Azure Cosmos DB 账号 URI,与 `endpoint` 视为同义配置。
+
+### endpoint [string]
+
+Azure Cosmos DB 账号 endpoint,例如
`https://example-account.documents.azure.com:443/`。
+
+### key [string]
+
+Azure Cosmos DB 账号 key。
+
+### primary_key [string]
+
+Azure Cosmos DB 主账号 key。
+
+### secondary_key [string]
+
+Azure Cosmos DB 次账号 key。
+
+### primary_connection_string [string]
+
+Azure Cosmos DB 主连接字符串。连接器可从中解析 `AccountEndpoint` 与 `AccountKey`。
+
+### secondary_connection_string [string]
+
+Azure Cosmos DB 次连接字符串。连接器可从中解析 `AccountEndpoint` 与 `AccountKey`。
+
+### database [string]
+
+目标数据库名。数据库必须已存在。
+
+### container [string]
+
+目标容器名。容器必须已存在。
+
+### schema [config]
+
+Cosmos DB 存储 JSON 文档,不会自动推断 SeaTunnel schema。必须显式配置 `schema.fields`。更多信息请参考
[Schema 特性](../../introduction/concepts/schema-feature.md)。
+
+示例:
+
+```hocon
+schema = {
+ fields {
+ id = string
+ user_id = string
+ amount = double
+ created_at = timestamp
+ labels = "map<string,string>"
+ }
+}
+```
+
+### query [string]
+
+用于有界批处理读取的 Cosmos SQL。可用于过滤和字段投影,例如 `SELECT c.id, c.name FROM c WHERE c.score
> 10`。
+
+### max_item_count [int]
+
+SDK 迭代查询结果时的首选分页大小。执行 checkpoint 恢复时,连接器会持久化当前分页查询的 SDK continuation
token,并从上一个完成的分页继续读取。
+
+### common-options
+
+源插件通用参数,请参考 [Source Common
Options](../common-options/source-common-options.md)。
+
+### 提示
+
+> 1. 必须至少提供
`uri`、`endpoint`、`primary_connection_string`、`secondary_connection_string`
之一,以及 `key`、`primary_key`、`secondary_key` 之一,或包含 account key 的连接字符串。<br/>
+> 2. V1 仅使用单个 split。提高 source 并行度不会按物理分区并行读取 Cosmos 数据。<br/>
+> 3. 可在 `query` 中使用 Cosmos SQL 做过滤和投影。连接器不提供单独的列投影特性。<br/>
+> 4. 连接器只读取已有容器,不会创建数据库、容器或索引。
+> 5. Checkpoint 恢复基于 Cosmos 查询分页边界,而不是单行边界。Change feed 读取仍不在范围内。
+
+## 如何创建 Azure Cosmos DB 数据同步作业
+
+以下示例演示如何创建批处理作业,从 Azure Cosmos DB 读取数据并打印到本地客户端:
+
+```bash
+# 设置要执行的任务的基本配置
+env {
+ parallelism = 1
+ job.mode = "BATCH"
+}
+
+source {
+ AzureCosmosDB {
+ uri = "https://example-account.documents.azure.com:443/"
+ primary_key = "<cosmos-account-key>"
+ database = "app-db"
+ container = "orders"
+ query = "SELECT c.id, c.user_id, c.amount, c.created_at, c.labels FROM c"
+ max_item_count = 200
+ schema = {
+ fields {
+ id = string
+ user_id = string
+ amount = double
+ created_at = timestamp
+ labels = "map<string,string>"
+ }
+ }
+ }
+}
+
+sink {
+ Console {
+ parallelism = 1
+ }
+}
+```
+
+### 带过滤条件与较小分页大小的查询
+
+```bash
+env {
+ parallelism = 1
+ job.mode = "BATCH"
+}
+
+source {
+ AzureCosmosDB {
+ endpoint = "https://example-account.documents.azure.com:443/"
+ key = "<cosmos-account-key>"
+ database = "app-db"
+ container = "orders"
+ query = "SELECT c.id, c.user_id, c.amount FROM c WHERE c.amount > 100"
+ max_item_count = 50
+ schema = {
+ fields {
+ id = string
+ user_id = string
+ amount = double
+ }
+ }
+ }
+}
+
+sink {
+ Console {}
+}
+```
+
+## 变更日志
+
+<ChangeLog />
diff --git a/plugin-mapping.properties b/plugin-mapping.properties
index 4ef59376d1..ea41a4b954 100644
--- a/plugin-mapping.properties
+++ b/plugin-mapping.properties
@@ -84,6 +84,7 @@ seatunnel.source.S3File = connector-file-s3
seatunnel.sink.S3File = connector-file-s3
seatunnel.source.AmazonDynamodb = connector-amazondynamodb
seatunnel.sink.AmazonDynamodb = connector-amazondynamodb
+seatunnel.source.AzureCosmosDB = connector-azurecosmosdb
seatunnel.source.Cassandra = connector-cassandra
seatunnel.sink.Cassandra = connector-cassandra
seatunnel.source.Salesforce = connector-salesforce
diff --git a/seatunnel-connectors-v2/connector-azurecosmosdb/pom.xml
b/seatunnel-connectors-v2/connector-azurecosmosdb/pom.xml
new file mode 100644
index 0000000000..5b3a4949b3
--- /dev/null
+++ b/seatunnel-connectors-v2/connector-azurecosmosdb/pom.xml
@@ -0,0 +1,45 @@
+<?xml version="1.0" encoding="UTF-8"?>
+<!--
+
+ 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.
+
+-->
+<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
+ xsi:schemaLocation="http://maven.apache.org/POM/4.0.0
http://maven.apache.org/xsd/maven-4.0.0.xsd">
+ <modelVersion>4.0.0</modelVersion>
+ <parent>
+ <groupId>org.apache.seatunnel</groupId>
+ <artifactId>seatunnel-connectors-v2</artifactId>
+ <version>${revision}</version>
+ </parent>
+
+ <artifactId>connector-azurecosmosdb</artifactId>
+ <name>SeaTunnel : Connectors V2 : Azure Cosmos DB</name>
+
+ <dependencies>
+ <dependency>
+ <groupId>org.apache.seatunnel</groupId>
+ <artifactId>connector-common</artifactId>
+ <version>${project.version}</version>
+ </dependency>
+ <dependency>
+ <groupId>com.azure</groupId>
+ <artifactId>azure-cosmos</artifactId>
+ <version>4.63.0</version>
+ </dependency>
+ </dependencies>
+
+</project>
diff --git
a/seatunnel-connectors-v2/connector-azurecosmosdb/src/main/java/org/apache/seatunnel/connectors/seatunnel/azurecosmosdb/config/AzureCosmosDBConfig.java
b/seatunnel-connectors-v2/connector-azurecosmosdb/src/main/java/org/apache/seatunnel/connectors/seatunnel/azurecosmosdb/config/AzureCosmosDBConfig.java
new file mode 100644
index 0000000000..238936c3d3
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-azurecosmosdb/src/main/java/org/apache/seatunnel/connectors/seatunnel/azurecosmosdb/config/AzureCosmosDBConfig.java
@@ -0,0 +1,181 @@
+/*
+ * 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 org.apache.seatunnel.connectors.seatunnel.azurecosmosdb.config;
+
+import org.apache.seatunnel.shade.com.typesafe.config.Config;
+
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.options.ConnectorCommonOptions;
+
+import java.io.Serializable;
+import java.util.HashMap;
+import java.util.Map;
+
+public class AzureCosmosDBConfig implements Serializable {
+
+ private final String uri;
+ private final String endpoint;
+ private final String key;
+ private final String primaryKey;
+ private final String secondaryKey;
+ private final String primaryConnectionString;
+ private final String secondaryConnectionString;
+ private final String database;
+ private final String container;
+ private final String query;
+ private final int maxItemCount;
+ private final Config schema;
+
+ public AzureCosmosDBConfig(ReadonlyConfig config) {
+ this.uri =
config.getOptional(AzureCosmosDBSourceOptions.URI).orElse(null);
+ this.endpoint =
config.getOptional(AzureCosmosDBSourceOptions.ENDPOINT).orElse(null);
+ this.key =
config.getOptional(AzureCosmosDBSourceOptions.KEY).orElse(null);
+ this.primaryKey =
config.getOptional(AzureCosmosDBSourceOptions.PRIMARY_KEY).orElse(null);
+ this.secondaryKey =
+
config.getOptional(AzureCosmosDBSourceOptions.SECONDARY_KEY).orElse(null);
+ this.primaryConnectionString =
+
config.getOptional(AzureCosmosDBSourceOptions.PRIMARY_CONNECTION_STRING)
+ .orElse(null);
+ this.secondaryConnectionString =
+
config.getOptional(AzureCosmosDBSourceOptions.SECONDARY_CONNECTION_STRING)
+ .orElse(null);
+ this.database = config.get(AzureCosmosDBSourceOptions.DATABASE);
+ this.container = config.get(AzureCosmosDBSourceOptions.CONTAINER);
+ this.query = config.get(AzureCosmosDBSourceOptions.QUERY);
+ this.maxItemCount =
config.get(AzureCosmosDBSourceOptions.MAX_ITEM_COUNT);
+ this.schema =
+ config.getOptional(ConnectorCommonOptions.SCHEMA)
+ .map(ReadonlyConfig::fromMap)
+ .map(ReadonlyConfig::toConfig)
+ .orElse(null);
+
+ if (getResolvedEndpoint() == null) {
+ throw new IllegalArgumentException(
+ "AzureCosmosDB requires uri, endpoint, or connection
string to resolve the endpoint");
+ }
+ if (getResolvedKey() == null) {
+ throw new IllegalArgumentException(
+ "AzureCosmosDB requires key, primary_key, secondary_key,
or a connection string to resolve the key");
+ }
+ }
+
+ public String getResolvedEndpoint() {
+ String resolvedEndpoint = firstNonBlank(uri, endpoint);
+ if (resolvedEndpoint != null) {
+ return resolvedEndpoint;
+ }
+
+ return firstNonBlank(
+ parseConnectionString(primaryConnectionString).get("endpoint"),
+
parseConnectionString(secondaryConnectionString).get("endpoint"));
+ }
+
+ public String getResolvedKey() {
+ String resolvedKey = firstNonBlank(key, primaryKey, secondaryKey);
+ if (resolvedKey != null) {
+ return resolvedKey;
+ }
+
+ return firstNonBlank(
+ parseConnectionString(primaryConnectionString).get("key"),
+ parseConnectionString(secondaryConnectionString).get("key"));
+ }
+
+ private static String firstNonBlank(String... values) {
+ for (String value : values) {
+ if (value != null && !value.trim().isEmpty()) {
+ return value.trim();
+ }
+ }
+ return null;
+ }
+
+ private static Map<String, String> parseConnectionString(String
connectionString) {
+ Map<String, String> values = new HashMap<>();
+ if (connectionString == null || connectionString.trim().isEmpty()) {
+ return values;
+ }
+ String[] parts = connectionString.split(";");
+ for (String part : parts) {
+ String trimmed = part.trim();
+ if (trimmed.isEmpty()) {
+ continue;
+ }
+ int equalsIndex = trimmed.indexOf('=');
+ if (equalsIndex <= 0 || equalsIndex == trimmed.length() - 1) {
+ continue;
+ }
+ String key = trimmed.substring(0,
equalsIndex).trim().toLowerCase();
+ String value = trimmed.substring(equalsIndex + 1).trim();
+ if ("accountendpoint".equals(key)) {
+ values.put("endpoint", value);
+ } else if ("accountkey".equals(key)) {
+ values.put("key", value);
+ }
+ }
+ return values;
+ }
+
+ public String getUri() {
+ return uri;
+ }
+
+ public String getEndpoint() {
+ return endpoint;
+ }
+
+ public String getKey() {
+ return key;
+ }
+
+ public String getPrimaryKey() {
+ return primaryKey;
+ }
+
+ public String getSecondaryKey() {
+ return secondaryKey;
+ }
+
+ public String getPrimaryConnectionString() {
+ return primaryConnectionString;
+ }
+
+ public String getSecondaryConnectionString() {
+ return secondaryConnectionString;
+ }
+
+ public String getDatabase() {
+ return database;
+ }
+
+ public String getContainer() {
+ return container;
+ }
+
+ public String getQuery() {
+ return query;
+ }
+
+ public int getMaxItemCount() {
+ return maxItemCount;
+ }
+
+ public Config getSchema() {
+ return schema;
+ }
+}
diff --git
a/seatunnel-connectors-v2/connector-azurecosmosdb/src/main/java/org/apache/seatunnel/connectors/seatunnel/azurecosmosdb/config/AzureCosmosDBSourceOptions.java
b/seatunnel-connectors-v2/connector-azurecosmosdb/src/main/java/org/apache/seatunnel/connectors/seatunnel/azurecosmosdb/config/AzureCosmosDBSourceOptions.java
new file mode 100644
index 0000000000..cbfa19150d
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-azurecosmosdb/src/main/java/org/apache/seatunnel/connectors/seatunnel/azurecosmosdb/config/AzureCosmosDBSourceOptions.java
@@ -0,0 +1,93 @@
+/*
+ * 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 org.apache.seatunnel.connectors.seatunnel.azurecosmosdb.config;
+
+import org.apache.seatunnel.api.configuration.Option;
+import org.apache.seatunnel.api.configuration.Options;
+import org.apache.seatunnel.api.options.ConnectorCommonOptions;
+
+import java.io.Serializable;
+
+public class AzureCosmosDBSourceOptions extends ConnectorCommonOptions
implements Serializable {
+
+ public static final Option<String> URI =
+ Options.key("uri")
+ .stringType()
+ .noDefaultValue()
+ .withDescription("Azure Cosmos DB account URI");
+
+ public static final Option<String> ENDPOINT =
+ Options.key("endpoint")
+ .stringType()
+ .noDefaultValue()
+ .withDescription("Azure Cosmos DB account endpoint");
+
+ public static final Option<String> KEY =
+ Options.key("key")
+ .stringType()
+ .noDefaultValue()
+ .withDescription("Azure Cosmos DB account key");
+
+ public static final Option<String> PRIMARY_KEY =
+ Options.key("primary_key")
+ .stringType()
+ .noDefaultValue()
+ .withDescription("Azure Cosmos DB primary account key");
+
+ public static final Option<String> SECONDARY_KEY =
+ Options.key("secondary_key")
+ .stringType()
+ .noDefaultValue()
+ .withDescription("Azure Cosmos DB secondary account key");
+
+ public static final Option<String> PRIMARY_CONNECTION_STRING =
+ Options.key("primary_connection_string")
+ .stringType()
+ .noDefaultValue()
+ .withDescription("Azure Cosmos DB primary connection
string");
+
+ public static final Option<String> SECONDARY_CONNECTION_STRING =
+ Options.key("secondary_connection_string")
+ .stringType()
+ .noDefaultValue()
+ .withDescription("Azure Cosmos DB secondary connection
string");
+
+ public static final Option<String> DATABASE =
+ Options.key("database")
+ .stringType()
+ .noDefaultValue()
+ .withDescription("Azure Cosmos DB database name");
+
+ public static final Option<String> CONTAINER =
+ Options.key("container")
+ .stringType()
+ .noDefaultValue()
+ .withDescription("Azure Cosmos DB container name");
+
+ public static final Option<String> QUERY =
+ Options.key("query")
+ .stringType()
+ .defaultValue("SELECT * FROM c")
+ .withDescription("Cosmos SQL query used to read source
data");
+
+ public static final Option<Integer> MAX_ITEM_COUNT =
+ Options.key("max_item_count")
+ .intType()
+ .defaultValue(100)
+ .withDescription("Max item count per Cosmos query page");
+}
diff --git
a/seatunnel-connectors-v2/connector-azurecosmosdb/src/main/java/org/apache/seatunnel/connectors/seatunnel/azurecosmosdb/serialize/CosmosItemDeserializer.java
b/seatunnel-connectors-v2/connector-azurecosmosdb/src/main/java/org/apache/seatunnel/connectors/seatunnel/azurecosmosdb/serialize/CosmosItemDeserializer.java
new file mode 100644
index 0000000000..718f42f5b8
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-azurecosmosdb/src/main/java/org/apache/seatunnel/connectors/seatunnel/azurecosmosdb/serialize/CosmosItemDeserializer.java
@@ -0,0 +1,154 @@
+/*
+ * 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 org.apache.seatunnel.connectors.seatunnel.azurecosmosdb.serialize;
+
+import org.apache.seatunnel.shade.com.fasterxml.jackson.databind.JsonNode;
+import org.apache.seatunnel.shade.com.fasterxml.jackson.databind.ObjectMapper;
+
+import org.apache.seatunnel.api.table.type.ArrayType;
+import org.apache.seatunnel.api.table.type.MapType;
+import org.apache.seatunnel.api.table.type.SeaTunnelDataType;
+import org.apache.seatunnel.api.table.type.SeaTunnelRow;
+import org.apache.seatunnel.api.table.type.SeaTunnelRowType;
+import org.apache.seatunnel.common.exception.CommonError;
+
+import java.lang.reflect.Array;
+import java.math.BigDecimal;
+import java.time.LocalDate;
+import java.time.LocalDateTime;
+import java.time.LocalTime;
+import java.util.HashMap;
+import java.util.Iterator;
+import java.util.Map;
+
+public class CosmosItemDeserializer {
+
+ private static final ObjectMapper OBJECT_MAPPER = new ObjectMapper();
+
+ private final SeaTunnelRowType rowType;
+
+ public CosmosItemDeserializer(SeaTunnelRowType rowType) {
+ this.rowType = rowType;
+ }
+
+ public SeaTunnelRow deserialize(Object item) {
+ JsonNode root = OBJECT_MAPPER.valueToTree(item);
+ SeaTunnelDataType<?>[] fieldTypes = rowType.getFieldTypes();
+ String[] fieldNames = rowType.getFieldNames();
+ Object[] fields = new Object[fieldNames.length];
+
+ for (int i = 0; i < fieldNames.length; i++) {
+ fields[i] = convert(fieldNames[i], fieldTypes[i],
root.get(fieldNames[i]));
+ }
+ return new SeaTunnelRow(fields);
+ }
+
+ private Object convert(String field, SeaTunnelDataType<?> type, JsonNode
node) {
+ if (node == null || node.isNull()) {
+ return null;
+ }
+
+ switch (type.getSqlType()) {
+ case BOOLEAN:
+ return node.asBoolean();
+ case TINYINT:
+ return (byte) node.asInt();
+ case SMALLINT:
+ return (short) node.asInt();
+ case INT:
+ return node.asInt();
+ case BIGINT:
+ return node.asLong();
+ case FLOAT:
+ return (float) node.asDouble();
+ case DOUBLE:
+ return node.asDouble();
+ case DECIMAL:
+ return new BigDecimal(node.asText());
+ case STRING:
+ return node.isTextual() ? node.asText() : node.toString();
+ case DATE:
+ return LocalDate.parse(node.asText());
+ case TIME:
+ return LocalTime.parse(node.asText());
+ case TIMESTAMP:
+ return LocalDateTime.parse(node.asText());
+ case BYTES:
+ try {
+ return node.binaryValue();
+ } catch (Exception e) {
+ throw CommonError.convertToSeaTunnelTypeError(
+ "AzureCosmosDB", type.getSqlType().toString(),
field);
+ }
+ case MAP:
+ return convertMap(field, (MapType<?, ?>) type, node);
+ case ARRAY:
+ return convertArray(field, (ArrayType<?, ?>) type, node);
+ case ROW:
+ return convertRow(field, (SeaTunnelRowType) type, node);
+ default:
+ throw CommonError.convertToSeaTunnelTypeError(
+ "AzureCosmosDB", type.getSqlType().toString(), field);
+ }
+ }
+
+ private Map<Object, Object> convertMap(String field, MapType<?, ?>
mapType, JsonNode node) {
+ if (!node.isObject()) {
+ throw CommonError.convertToSeaTunnelTypeError(
+ "AzureCosmosDB", mapType.getSqlType().toString(), field);
+ }
+
+ Map<Object, Object> values = new HashMap<>();
+ Iterator<Map.Entry<String, JsonNode>> fields = node.fields();
+ while (fields.hasNext()) {
+ Map.Entry<String, JsonNode> entry = fields.next();
+ Object key =
+ convert(field, mapType.getKeyType(),
OBJECT_MAPPER.valueToTree(entry.getKey()));
+ Object value = convert(field, mapType.getValueType(),
entry.getValue());
+ values.put(key, value);
+ }
+ return values;
+ }
+
+ private Object convertArray(String field, ArrayType<?, ?> arrayType,
JsonNode node) {
+ if (!node.isArray()) {
+ throw CommonError.convertToSeaTunnelTypeError(
+ "AzureCosmosDB", arrayType.getSqlType().toString(), field);
+ }
+
+ Object array =
Array.newInstance(arrayType.getElementType().getTypeClass(), node.size());
+ for (int i = 0; i < node.size(); i++) {
+ Array.set(array, i, convert(field, arrayType.getElementType(),
node.get(i)));
+ }
+ return array;
+ }
+
+ private SeaTunnelRow convertRow(String field, SeaTunnelRowType rowType,
JsonNode node) {
+ if (!node.isObject()) {
+ throw CommonError.convertToSeaTunnelTypeError(
+ "AzureCosmosDB", rowType.getSqlType().toString(), field);
+ }
+
+ Object[] fields = new Object[rowType.getTotalFields()];
+ for (int i = 0; i < rowType.getTotalFields(); i++) {
+ String fieldName = rowType.getFieldName(i);
+ fields[i] = convert(fieldName, rowType.getFieldType(i),
node.get(fieldName));
+ }
+ return new SeaTunnelRow(fields);
+ }
+}
diff --git
a/seatunnel-connectors-v2/connector-azurecosmosdb/src/main/java/org/apache/seatunnel/connectors/seatunnel/azurecosmosdb/source/AzureCosmosDBSource.java
b/seatunnel-connectors-v2/connector-azurecosmosdb/src/main/java/org/apache/seatunnel/connectors/seatunnel/azurecosmosdb/source/AzureCosmosDBSource.java
new file mode 100644
index 0000000000..77133f0c8c
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-azurecosmosdb/src/main/java/org/apache/seatunnel/connectors/seatunnel/azurecosmosdb/source/AzureCosmosDBSource.java
@@ -0,0 +1,81 @@
+/*
+ * 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 org.apache.seatunnel.connectors.seatunnel.azurecosmosdb.source;
+
+import org.apache.seatunnel.api.source.Boundedness;
+import org.apache.seatunnel.api.source.SeaTunnelSource;
+import org.apache.seatunnel.api.source.SourceReader;
+import org.apache.seatunnel.api.source.SourceSplitEnumerator;
+import org.apache.seatunnel.api.table.catalog.CatalogTable;
+import org.apache.seatunnel.api.table.type.SeaTunnelRow;
+import
org.apache.seatunnel.connectors.seatunnel.azurecosmosdb.config.AzureCosmosDBConfig;
+
+import java.util.Collections;
+import java.util.List;
+
+public class AzureCosmosDBSource
+ implements SeaTunnelSource<
+ SeaTunnelRow, AzureCosmosDBSourceSplit,
AzureCosmosDBSourceState> {
+
+ private final AzureCosmosDBConfig config;
+ private final CatalogTable catalogTable;
+
+ public AzureCosmosDBSource(AzureCosmosDBConfig config, CatalogTable
catalogTable) {
+ this.config = config;
+ this.catalogTable = catalogTable;
+ }
+
+ @Override
+ public String getPluginName() {
+ return "AzureCosmosDB";
+ }
+
+ @Override
+ public Boundedness getBoundedness() {
+ return Boundedness.BOUNDED;
+ }
+
+ @Override
+ public List<CatalogTable> getProducedCatalogTables() {
+ return Collections.singletonList(catalogTable);
+ }
+
+ @Override
+ public SourceSplitEnumerator<AzureCosmosDBSourceSplit,
AzureCosmosDBSourceState>
+ createEnumerator(
+ SourceSplitEnumerator.Context<AzureCosmosDBSourceSplit>
enumeratorContext)
+ throws Exception {
+ return new AzureCosmosDBSourceSplitEnumerator(enumeratorContext, null);
+ }
+
+ @Override
+ public SourceSplitEnumerator<AzureCosmosDBSourceSplit,
AzureCosmosDBSourceState>
+ restoreEnumerator(
+ SourceSplitEnumerator.Context<AzureCosmosDBSourceSplit>
enumeratorContext,
+ AzureCosmosDBSourceState checkpointState)
+ throws Exception {
+ return new AzureCosmosDBSourceSplitEnumerator(enumeratorContext,
checkpointState);
+ }
+
+ @Override
+ public SourceReader<SeaTunnelRow, AzureCosmosDBSourceSplit> createReader(
+ SourceReader.Context readerContext) throws Exception {
+ return new AzureCosmosDBSourceReader(
+ readerContext, config, catalogTable.getSeaTunnelRowType());
+ }
+}
diff --git
a/seatunnel-connectors-v2/connector-azurecosmosdb/src/main/java/org/apache/seatunnel/connectors/seatunnel/azurecosmosdb/source/AzureCosmosDBSourceFactory.java
b/seatunnel-connectors-v2/connector-azurecosmosdb/src/main/java/org/apache/seatunnel/connectors/seatunnel/azurecosmosdb/source/AzureCosmosDBSourceFactory.java
new file mode 100644
index 0000000000..8f209bf490
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-azurecosmosdb/src/main/java/org/apache/seatunnel/connectors/seatunnel/azurecosmosdb/source/AzureCosmosDBSourceFactory.java
@@ -0,0 +1,86 @@
+/*
+ * 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 org.apache.seatunnel.connectors.seatunnel.azurecosmosdb.source;
+
+import org.apache.seatunnel.api.configuration.util.OptionRule;
+import org.apache.seatunnel.api.source.SeaTunnelSource;
+import org.apache.seatunnel.api.source.SourceSplit;
+import org.apache.seatunnel.api.table.catalog.CatalogTableUtil;
+import org.apache.seatunnel.api.table.connector.TableSource;
+import org.apache.seatunnel.api.table.factory.Factory;
+import org.apache.seatunnel.api.table.factory.TableSourceFactory;
+import org.apache.seatunnel.api.table.factory.TableSourceFactoryContext;
+import
org.apache.seatunnel.connectors.seatunnel.azurecosmosdb.config.AzureCosmosDBConfig;
+
+import com.google.auto.service.AutoService;
+
+import java.io.Serializable;
+
+import static org.apache.seatunnel.api.options.ConnectorCommonOptions.SCHEMA;
+import static
org.apache.seatunnel.connectors.seatunnel.azurecosmosdb.config.AzureCosmosDBSourceOptions.CONTAINER;
+import static
org.apache.seatunnel.connectors.seatunnel.azurecosmosdb.config.AzureCosmosDBSourceOptions.DATABASE;
+import static
org.apache.seatunnel.connectors.seatunnel.azurecosmosdb.config.AzureCosmosDBSourceOptions.ENDPOINT;
+import static
org.apache.seatunnel.connectors.seatunnel.azurecosmosdb.config.AzureCosmosDBSourceOptions.KEY;
+import static
org.apache.seatunnel.connectors.seatunnel.azurecosmosdb.config.AzureCosmosDBSourceOptions.MAX_ITEM_COUNT;
+import static
org.apache.seatunnel.connectors.seatunnel.azurecosmosdb.config.AzureCosmosDBSourceOptions.PRIMARY_CONNECTION_STRING;
+import static
org.apache.seatunnel.connectors.seatunnel.azurecosmosdb.config.AzureCosmosDBSourceOptions.PRIMARY_KEY;
+import static
org.apache.seatunnel.connectors.seatunnel.azurecosmosdb.config.AzureCosmosDBSourceOptions.QUERY;
+import static
org.apache.seatunnel.connectors.seatunnel.azurecosmosdb.config.AzureCosmosDBSourceOptions.SECONDARY_CONNECTION_STRING;
+import static
org.apache.seatunnel.connectors.seatunnel.azurecosmosdb.config.AzureCosmosDBSourceOptions.SECONDARY_KEY;
+import static
org.apache.seatunnel.connectors.seatunnel.azurecosmosdb.config.AzureCosmosDBSourceOptions.URI;
+
+@AutoService(Factory.class)
+public class AzureCosmosDBSourceFactory implements TableSourceFactory {
+
+ @Override
+ public String factoryIdentifier() {
+ return "AzureCosmosDB";
+ }
+
+ @Override
+ public OptionRule optionRule() {
+ return OptionRule.builder()
+ .required(DATABASE, CONTAINER, SCHEMA)
+ .optional(
+ URI,
+ ENDPOINT,
+ KEY,
+ PRIMARY_KEY,
+ SECONDARY_KEY,
+ PRIMARY_CONNECTION_STRING,
+ SECONDARY_CONNECTION_STRING,
+ QUERY,
+ MAX_ITEM_COUNT)
+ .build();
+ }
+
+ @Override
+ public <T, SplitT extends SourceSplit, StateT extends Serializable>
+ TableSource<T, SplitT, StateT>
createSource(TableSourceFactoryContext context) {
+ return () ->
+ (SeaTunnelSource<T, SplitT, StateT>)
+ new AzureCosmosDBSource(
+ new AzureCosmosDBConfig(context.getOptions()),
+
CatalogTableUtil.buildWithConfig(context.getOptions()));
+ }
+
+ @Override
+ public Class<? extends SeaTunnelSource> getSourceClass() {
+ return AzureCosmosDBSource.class;
+ }
+}
diff --git
a/seatunnel-connectors-v2/connector-azurecosmosdb/src/main/java/org/apache/seatunnel/connectors/seatunnel/azurecosmosdb/source/AzureCosmosDBSourceReader.java
b/seatunnel-connectors-v2/connector-azurecosmosdb/src/main/java/org/apache/seatunnel/connectors/seatunnel/azurecosmosdb/source/AzureCosmosDBSourceReader.java
new file mode 100644
index 0000000000..b6cb98a921
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-azurecosmosdb/src/main/java/org/apache/seatunnel/connectors/seatunnel/azurecosmosdb/source/AzureCosmosDBSourceReader.java
@@ -0,0 +1,242 @@
+/*
+ * 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 org.apache.seatunnel.connectors.seatunnel.azurecosmosdb.source;
+
+import org.apache.seatunnel.api.source.Collector;
+import org.apache.seatunnel.api.source.SourceReader;
+import org.apache.seatunnel.api.table.type.SeaTunnelRow;
+import org.apache.seatunnel.api.table.type.SeaTunnelRowType;
+import
org.apache.seatunnel.connectors.seatunnel.azurecosmosdb.config.AzureCosmosDBConfig;
+import
org.apache.seatunnel.connectors.seatunnel.azurecosmosdb.serialize.CosmosItemDeserializer;
+
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import com.azure.cosmos.CosmosClient;
+import com.azure.cosmos.CosmosClientBuilder;
+import com.azure.cosmos.CosmosContainer;
+import com.azure.cosmos.models.CosmosQueryRequestOptions;
+import com.azure.cosmos.models.FeedResponse;
+
+import java.util.ArrayList;
+import java.util.Iterator;
+import java.util.List;
+import java.util.Objects;
+import java.util.Queue;
+import java.util.concurrent.ConcurrentLinkedDeque;
+
+public class AzureCosmosDBSourceReader
+ implements SourceReader<SeaTunnelRow, AzureCosmosDBSourceSplit> {
+
+ private static final Logger LOG =
LoggerFactory.getLogger(AzureCosmosDBSourceReader.class);
+
+ private final SourceReader.Context context;
+ private final AzureCosmosDBConfig config;
+ private final CosmosItemDeserializer deserializer;
+
+ private final Queue<AzureCosmosDBSourceSplit> pendingSplits = new
ConcurrentLinkedDeque<>();
+
+ private CosmosClient client;
+ private CosmosContainer container;
+ private AzureCosmosDBSourceSplit currentSplit;
+
+ private volatile boolean noMoreSplit;
+ private volatile boolean finished;
+
+ public AzureCosmosDBSourceReader(
+ SourceReader.Context context, AzureCosmosDBConfig config,
SeaTunnelRowType rowType) {
+ this.context = context;
+ this.config = config;
+ this.deserializer = new CosmosItemDeserializer(rowType);
+ }
+
+ @Override
+ public void open() {
+ try {
+ this.client =
+ new CosmosClientBuilder()
+ .endpoint(config.getResolvedEndpoint())
+ .key(config.getResolvedKey())
+ .endpointDiscoveryEnabled(false)
+ .gatewayMode()
+ .buildClient();
+ this.container =
+
client.getDatabase(config.getDatabase()).getContainer(config.getContainer());
+ } catch (Exception e) {
+ throw new IllegalStateException(
+ String.format(
+ "Failed to open AzureCosmosDB source reader for
database [%s], container [%s]",
+ config.getDatabase(), config.getContainer()),
+ e);
+ }
+ }
+
+ @Override
+ public void close() {
+ if (client != null) {
+ client.close();
+ }
+ }
+
+ @Override
+ public void pollNext(Collector<SeaTunnelRow> output) {
+ AzureCosmosDBSourceSplit activeSplit;
+ String continuationToken;
+
+ synchronized (output.getCheckpointLock()) {
+ if (finished) {
+ return;
+ }
+
+ if (currentSplit == null) {
+ currentSplit = pendingSplits.poll();
+ }
+ if (currentSplit == null) {
+ if (noMoreSplit) {
+ finishReader();
+ }
+ return;
+ }
+
+ activeSplit = currentSplit;
+ continuationToken = activeSplit.getContinuationToken();
+ }
+
+ FeedResponse<Object> page = fetchPage(continuationToken);
+ List<SeaTunnelRow> rows = deserializePage(page);
+
+ synchronized (output.getCheckpointLock()) {
+ applyPage(activeSplit, page, rows, output);
+ finishReaderIfNoMoreWork();
+ }
+ }
+
+ @Override
+ public List<AzureCosmosDBSourceSplit> snapshotState(long checkpointId) {
+ List<AzureCosmosDBSourceSplit> state = new ArrayList<>();
+ pendingSplits.forEach(split -> state.add(split.copy()));
+ if (currentSplit != null) {
+ state.add(currentSplit.copy());
+ }
+ return state;
+ }
+
+ @Override
+ public void addSplits(List<AzureCosmosDBSourceSplit> splits) {
+ pendingSplits.addAll(splits);
+ }
+
+ @Override
+ public void handleNoMoreSplits() {
+ this.noMoreSplit = true;
+ }
+
+ @Override
+ public void notifyCheckpointComplete(long checkpointId) {
+ // no-op
+ }
+
+ private List<SeaTunnelRow> deserializePage(FeedResponse<Object> page) {
+ List<SeaTunnelRow> rows = new ArrayList<>();
+ if (page == null) {
+ return rows;
+ }
+
+ for (Object item : page.getResults()) {
+ if (Objects.nonNull(item)) {
+ rows.add(deserializer.deserialize(item));
+ }
+ }
+ return rows;
+ }
+
+ private void applyPage(
+ AzureCosmosDBSourceSplit activeSplit,
+ FeedResponse<Object> page,
+ List<SeaTunnelRow> rows,
+ Collector<SeaTunnelRow> output) {
+ if (finished || currentSplit != activeSplit) {
+ return;
+ }
+
+ if (page == null) {
+ finishCurrentSplit();
+ return;
+ }
+
+ rows.forEach(output::collect);
+
+ String continuationToken = page.getContinuationToken();
+ currentSplit.setContinuationToken(continuationToken);
+ if (isLastPage(continuationToken)) {
+ finishCurrentSplit();
+ }
+ }
+
+ FeedResponse<Object> fetchPage(String continuationToken) {
+ CosmosQueryRequestOptions queryOptions = new
CosmosQueryRequestOptions();
+
+ try {
+ Iterator<FeedResponse<Object>> pageIterator;
+ if (isLastPage(continuationToken)) {
+ pageIterator =
+ container
+ .<Object>queryItems(config.getQuery(),
queryOptions, Object.class)
+ .iterableByPage(config.getMaxItemCount())
+ .iterator();
+ } else {
+ pageIterator =
+ container
+ .<Object>queryItems(config.getQuery(),
queryOptions, Object.class)
+ .iterableByPage(continuationToken,
config.getMaxItemCount())
+ .iterator();
+ }
+ return pageIterator.hasNext() ? pageIterator.next() : null;
+ } catch (Exception e) {
+ throw new IllegalStateException(
+ String.format(
+ "Failed to read AzureCosmosDB data from database
[%s], container [%s] with query [%s]",
+ config.getDatabase(), config.getContainer(),
config.getQuery()),
+ e);
+ }
+ }
+
+ private static boolean isLastPage(String continuationToken) {
+ return continuationToken == null || continuationToken.isEmpty();
+ }
+
+ private void finishReader() {
+ context.signalNoMoreElement();
+ finished = true;
+ }
+
+ private void finishReaderIfNoMoreWork() {
+ if (currentSplit == null && pendingSplits.isEmpty() && noMoreSplit) {
+ finishReader();
+ }
+ }
+
+ private void finishCurrentSplit() {
+ LOG.info("AzureCosmosDB reader [{}] finished source scan",
context.getIndexOfSubtask());
+ currentSplit = null;
+ }
+
+ int getQueryPageSize() {
+ return config.getMaxItemCount();
+ }
+}
diff --git
a/seatunnel-connectors-v2/connector-azurecosmosdb/src/main/java/org/apache/seatunnel/connectors/seatunnel/azurecosmosdb/source/AzureCosmosDBSourceSplit.java
b/seatunnel-connectors-v2/connector-azurecosmosdb/src/main/java/org/apache/seatunnel/connectors/seatunnel/azurecosmosdb/source/AzureCosmosDBSourceSplit.java
new file mode 100644
index 0000000000..c34fb6e07c
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-azurecosmosdb/src/main/java/org/apache/seatunnel/connectors/seatunnel/azurecosmosdb/source/AzureCosmosDBSourceSplit.java
@@ -0,0 +1,58 @@
+/*
+ * 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 org.apache.seatunnel.connectors.seatunnel.azurecosmosdb.source;
+
+import org.apache.seatunnel.api.source.SourceSplit;
+
+public class AzureCosmosDBSourceSplit implements SourceSplit {
+
+ private static final long serialVersionUID = 2485413678354889739L;
+
+ private final Integer splitId;
+ private String continuationToken;
+
+ public AzureCosmosDBSourceSplit(Integer splitId) {
+ this(splitId, null);
+ }
+
+ public AzureCosmosDBSourceSplit(Integer splitId, String continuationToken)
{
+ this.splitId = splitId;
+ this.continuationToken = continuationToken;
+ }
+
+ public Integer getSplitId() {
+ return splitId;
+ }
+
+ public String getContinuationToken() {
+ return continuationToken;
+ }
+
+ public void setContinuationToken(String continuationToken) {
+ this.continuationToken = continuationToken;
+ }
+
+ @Override
+ public String splitId() {
+ return splitId.toString();
+ }
+
+ public AzureCosmosDBSourceSplit copy() {
+ return new AzureCosmosDBSourceSplit(splitId, continuationToken);
+ }
+}
diff --git
a/seatunnel-connectors-v2/connector-azurecosmosdb/src/main/java/org/apache/seatunnel/connectors/seatunnel/azurecosmosdb/source/AzureCosmosDBSourceSplitEnumerator.java
b/seatunnel-connectors-v2/connector-azurecosmosdb/src/main/java/org/apache/seatunnel/connectors/seatunnel/azurecosmosdb/source/AzureCosmosDBSourceSplitEnumerator.java
new file mode 100644
index 0000000000..edaa22126c
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-azurecosmosdb/src/main/java/org/apache/seatunnel/connectors/seatunnel/azurecosmosdb/source/AzureCosmosDBSourceSplitEnumerator.java
@@ -0,0 +1,145 @@
+/*
+ * 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 org.apache.seatunnel.connectors.seatunnel.azurecosmosdb.source;
+
+import org.apache.seatunnel.api.source.SourceSplitEnumerator;
+
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import java.io.IOException;
+import java.util.ArrayList;
+import java.util.Collection;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Set;
+
+public class AzureCosmosDBSourceSplitEnumerator
+ implements SourceSplitEnumerator<AzureCosmosDBSourceSplit,
AzureCosmosDBSourceState> {
+
+ private static final Logger LOG =
+ LoggerFactory.getLogger(AzureCosmosDBSourceSplitEnumerator.class);
+
+ private final SourceSplitEnumerator.Context<AzureCosmosDBSourceSplit>
enumeratorContext;
+ private final Map<Integer, List<AzureCosmosDBSourceSplit>> pendingSplits;
+ private final Object stateLock = new Object();
+
+ private volatile boolean shouldEnumerate;
+
+ public AzureCosmosDBSourceSplitEnumerator(
+ Context<AzureCosmosDBSourceSplit> enumeratorContext,
+ AzureCosmosDBSourceState sourceState) {
+ this.enumeratorContext = enumeratorContext;
+ this.pendingSplits = new HashMap<>();
+ this.shouldEnumerate = sourceState == null;
+ if (sourceState != null) {
+ this.shouldEnumerate = sourceState.isShouldEnumerate();
+ this.pendingSplits.putAll(sourceState.getPendingSplits());
+ }
+ }
+
+ @Override
+ public void open() {
+ // no-op
+ }
+
+ @Override
+ public void run() throws Exception {
+ Set<Integer> readers = enumeratorContext.registeredReaders();
+ if (shouldEnumerate) {
+ synchronized (stateLock) {
+ addPendingSplits(Collections.singletonList(new
AzureCosmosDBSourceSplit(0)));
+ shouldEnumerate = false;
+ }
+ assignSplit(readers);
+ }
+ }
+
+ @Override
+ public void close() throws IOException {
+ // no-op
+ }
+
+ @Override
+ public void addSplitsBack(List<AzureCosmosDBSourceSplit> splits, int
subtaskId) {
+ if (!splits.isEmpty()) {
+ addPendingSplits(splits);
+ assignSplit(Collections.singleton(subtaskId));
+ enumeratorContext.signalNoMoreSplits(subtaskId);
+ }
+ }
+
+ @Override
+ public int currentUnassignedSplitSize() {
+ return pendingSplits.size();
+ }
+
+ @Override
+ public void handleSplitRequest(int subtaskId) {
+ // no-op
+ }
+
+ @Override
+ public void registerReader(int subtaskId) {
+ if (!pendingSplits.isEmpty()) {
+ assignSplit(Collections.singleton(subtaskId));
+ }
+ }
+
+ @Override
+ public AzureCosmosDBSourceState snapshotState(long checkpointId) throws
Exception {
+ synchronized (stateLock) {
+ return new AzureCosmosDBSourceState(shouldEnumerate,
pendingSplits);
+ }
+ }
+
+ @Override
+ public void notifyCheckpointComplete(long checkpointId) throws Exception {
+ // no-op
+ }
+
+ private void addPendingSplits(Collection<AzureCosmosDBSourceSplit> splits)
{
+ int readerCount = enumeratorContext.currentParallelism();
+ for (AzureCosmosDBSourceSplit split : splits) {
+ int ownerReader = getSplitOwner(split.getSplitId(), readerCount);
+ pendingSplits.computeIfAbsent(ownerReader, id -> new
ArrayList<>()).add(split);
+ }
+ }
+
+ private void assignSplit(Set<Integer> readers) {
+ for (int reader : readers) {
+ List<AzureCosmosDBSourceSplit> assignment =
pendingSplits.remove(reader);
+ if (assignment != null && !assignment.isEmpty()) {
+ LOG.info("Assign splits {} to reader {}", assignment, reader);
+ try {
+ enumeratorContext.assignSplit(reader, assignment);
+ } catch (Exception e) {
+ LOG.error("Failed to assign splits {} to reader {}",
assignment, reader, e);
+ pendingSplits.put(reader, assignment);
+ }
+ }
+ enumeratorContext.signalNoMoreSplits(reader);
+ }
+ }
+
+ private static int getSplitOwner(Integer splitId, int numReaders) {
+ return (splitId.hashCode() & Integer.MAX_VALUE) % numReaders;
+ }
+}
diff --git
a/seatunnel-connectors-v2/connector-azurecosmosdb/src/main/java/org/apache/seatunnel/connectors/seatunnel/azurecosmosdb/source/AzureCosmosDBSourceState.java
b/seatunnel-connectors-v2/connector-azurecosmosdb/src/main/java/org/apache/seatunnel/connectors/seatunnel/azurecosmosdb/source/AzureCosmosDBSourceState.java
new file mode 100644
index 0000000000..1b78efe110
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-azurecosmosdb/src/main/java/org/apache/seatunnel/connectors/seatunnel/azurecosmosdb/source/AzureCosmosDBSourceState.java
@@ -0,0 +1,44 @@
+/*
+ * 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 org.apache.seatunnel.connectors.seatunnel.azurecosmosdb.source;
+
+import java.io.Serializable;
+import java.util.List;
+import java.util.Map;
+
+public class AzureCosmosDBSourceState implements Serializable {
+
+ private static final long serialVersionUID = -4911932288415348743L;
+
+ private final boolean shouldEnumerate;
+ private final Map<Integer, List<AzureCosmosDBSourceSplit>> pendingSplits;
+
+ public AzureCosmosDBSourceState(
+ boolean shouldEnumerate, Map<Integer,
List<AzureCosmosDBSourceSplit>> pendingSplits) {
+ this.shouldEnumerate = shouldEnumerate;
+ this.pendingSplits = pendingSplits;
+ }
+
+ public boolean isShouldEnumerate() {
+ return shouldEnumerate;
+ }
+
+ public Map<Integer, List<AzureCosmosDBSourceSplit>> getPendingSplits() {
+ return pendingSplits;
+ }
+}
diff --git
a/seatunnel-connectors-v2/connector-azurecosmosdb/src/test/java/org/apache/seatunnel/connectors/seatunnel/azurecosmosdb/AzureCosmosDBSourceFactoryTest.java
b/seatunnel-connectors-v2/connector-azurecosmosdb/src/test/java/org/apache/seatunnel/connectors/seatunnel/azurecosmosdb/AzureCosmosDBSourceFactoryTest.java
new file mode 100644
index 0000000000..48fcefc47a
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-azurecosmosdb/src/test/java/org/apache/seatunnel/connectors/seatunnel/azurecosmosdb/AzureCosmosDBSourceFactoryTest.java
@@ -0,0 +1,39 @@
+/*
+ * 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 org.apache.seatunnel.connectors.seatunnel.azurecosmosdb;
+
+import org.apache.seatunnel.api.configuration.util.OptionRule;
+import
org.apache.seatunnel.connectors.seatunnel.azurecosmosdb.source.AzureCosmosDBSourceFactory;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+import static
org.apache.seatunnel.connectors.seatunnel.azurecosmosdb.config.AzureCosmosDBSourceOptions.URI;
+
+public class AzureCosmosDBSourceFactoryTest {
+
+ @Test
+ public void testSourceFactory() {
+ AzureCosmosDBSourceFactory sourceFactory = new
AzureCosmosDBSourceFactory();
+ OptionRule optionRule = sourceFactory.optionRule();
+
+ Assertions.assertNotNull(optionRule);
+ Assertions.assertEquals("AzureCosmosDB",
sourceFactory.factoryIdentifier());
+ Assertions.assertTrue(optionRule.getOptionalOptions().contains(URI));
+ }
+}
diff --git
a/seatunnel-connectors-v2/connector-azurecosmosdb/src/test/java/org/apache/seatunnel/connectors/seatunnel/azurecosmosdb/CosmosItemDeserializerTest.java
b/seatunnel-connectors-v2/connector-azurecosmosdb/src/test/java/org/apache/seatunnel/connectors/seatunnel/azurecosmosdb/CosmosItemDeserializerTest.java
new file mode 100644
index 0000000000..f465929622
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-azurecosmosdb/src/test/java/org/apache/seatunnel/connectors/seatunnel/azurecosmosdb/CosmosItemDeserializerTest.java
@@ -0,0 +1,72 @@
+/*
+ * 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 org.apache.seatunnel.connectors.seatunnel.azurecosmosdb;
+
+import org.apache.seatunnel.api.table.type.ArrayType;
+import org.apache.seatunnel.api.table.type.BasicType;
+import org.apache.seatunnel.api.table.type.MapType;
+import org.apache.seatunnel.api.table.type.SeaTunnelDataType;
+import org.apache.seatunnel.api.table.type.SeaTunnelRow;
+import org.apache.seatunnel.api.table.type.SeaTunnelRowType;
+import
org.apache.seatunnel.connectors.seatunnel.azurecosmosdb.serialize.CosmosItemDeserializer;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+import java.util.HashMap;
+import java.util.Map;
+
+public class CosmosItemDeserializerTest {
+
+ @Test
+ public void testDeserialize() {
+ SeaTunnelRowType rowType =
+ new SeaTunnelRowType(
+ new String[] {"id", "name", "active", "score", "tags",
"labels"},
+ new SeaTunnelDataType[] {
+ BasicType.INT_TYPE,
+ BasicType.STRING_TYPE,
+ BasicType.BOOLEAN_TYPE,
+ BasicType.DOUBLE_TYPE,
+ ArrayType.of(BasicType.STRING_TYPE),
+ new MapType<>(BasicType.STRING_TYPE,
BasicType.STRING_TYPE)
+ });
+
+ CosmosItemDeserializer deserializer = new
CosmosItemDeserializer(rowType);
+ Map<String, Object> labels = new HashMap<>();
+ labels.put("region", "westus");
+ labels.put("team", "data");
+
+ Map<String, Object> doc = new HashMap<>();
+ doc.put("id", 7);
+ doc.put("name", "cosmos-user");
+ doc.put("active", true);
+ doc.put("score", 98.5);
+ doc.put("tags", new String[] {"alpha", "beta"});
+ doc.put("labels", labels);
+
+ SeaTunnelRow row = deserializer.deserialize(doc);
+
+ Assertions.assertEquals(7, row.getField(0));
+ Assertions.assertEquals("cosmos-user", row.getField(1));
+ Assertions.assertEquals(true, row.getField(2));
+ Assertions.assertEquals(98.5, row.getField(3));
+ Assertions.assertArrayEquals(new String[] {"alpha", "beta"},
(Object[]) row.getField(4));
+ Assertions.assertEquals(labels, row.getField(5));
+ }
+}
diff --git
a/seatunnel-connectors-v2/connector-azurecosmosdb/src/test/java/org/apache/seatunnel/connectors/seatunnel/azurecosmosdb/config/AzureCosmosDBConfigTest.java
b/seatunnel-connectors-v2/connector-azurecosmosdb/src/test/java/org/apache/seatunnel/connectors/seatunnel/azurecosmosdb/config/AzureCosmosDBConfigTest.java
new file mode 100644
index 0000000000..e73c92b68c
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-azurecosmosdb/src/test/java/org/apache/seatunnel/connectors/seatunnel/azurecosmosdb/config/AzureCosmosDBConfigTest.java
@@ -0,0 +1,69 @@
+/*
+ * 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 org.apache.seatunnel.connectors.seatunnel.azurecosmosdb.config;
+
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+import java.util.HashMap;
+import java.util.Map;
+
+public class AzureCosmosDBConfigTest {
+
+ @Test
+ public void testResolveFromUriAndPrimaryKey() {
+ AzureCosmosDBConfig config =
+ new
AzureCosmosDBConfig(ReadonlyConfig.fromMap(buildBaseOptions()));
+
+ Assertions.assertEquals(
+ "https://account.documents.azure.com:443/",
config.getResolvedEndpoint());
+ Assertions.assertEquals("primary-key", config.getResolvedKey());
+ }
+
+ @Test
+ public void testResolveFromConnectionString() {
+ Map<String, Object> options = buildBaseOptions();
+ options.remove("uri");
+ options.remove("primary_key");
+ options.put(
+ "primary_connection_string",
+
"AccountEndpoint=https://primary.documents.azure.com:443/;AccountKey=primary-connection-key;");
+
+ AzureCosmosDBConfig config = new
AzureCosmosDBConfig(ReadonlyConfig.fromMap(options));
+
+ Assertions.assertEquals(
+ "https://primary.documents.azure.com:443/",
config.getResolvedEndpoint());
+ Assertions.assertEquals("primary-connection-key",
config.getResolvedKey());
+ }
+
+ private Map<String, Object> buildBaseOptions() {
+ Map<String, Object> schema = new HashMap<>();
+ schema.put("fields", new HashMap<String, Object>());
+
+ Map<String, Object> options = new HashMap<>();
+ options.put("uri", "https://account.documents.azure.com:443/");
+ options.put("primary_key", "primary-key");
+ options.put("database", "test-db");
+ options.put("container", "test-container");
+ options.put("query", "SELECT * FROM c");
+ options.put("schema", schema);
+ return options;
+ }
+}
diff --git
a/seatunnel-connectors-v2/connector-azurecosmosdb/src/test/java/org/apache/seatunnel/connectors/seatunnel/azurecosmosdb/source/AzureCosmosDBSourceReaderTest.java
b/seatunnel-connectors-v2/connector-azurecosmosdb/src/test/java/org/apache/seatunnel/connectors/seatunnel/azurecosmosdb/source/AzureCosmosDBSourceReaderTest.java
new file mode 100644
index 0000000000..88884df13b
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-azurecosmosdb/src/test/java/org/apache/seatunnel/connectors/seatunnel/azurecosmosdb/source/AzureCosmosDBSourceReaderTest.java
@@ -0,0 +1,210 @@
+/*
+ * 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 org.apache.seatunnel.connectors.seatunnel.azurecosmosdb.source;
+
+import org.apache.seatunnel.api.common.metrics.MetricsContext;
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.event.EventListener;
+import org.apache.seatunnel.api.source.Boundedness;
+import org.apache.seatunnel.api.source.Collector;
+import org.apache.seatunnel.api.source.SourceEvent;
+import org.apache.seatunnel.api.source.SourceReader;
+import org.apache.seatunnel.api.table.type.BasicType;
+import org.apache.seatunnel.api.table.type.SeaTunnelDataType;
+import org.apache.seatunnel.api.table.type.SeaTunnelRow;
+import org.apache.seatunnel.api.table.type.SeaTunnelRowType;
+import
org.apache.seatunnel.connectors.seatunnel.azurecosmosdb.config.AzureCosmosDBConfig;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+import com.azure.cosmos.models.FeedResponse;
+
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.Map;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicReference;
+
+public class AzureCosmosDBSourceReaderTest {
+
+ @Test
+ public void testUsesConfiguredQueryPageSize() {
+ AzureCosmosDBSourceReader reader =
+ new AzureCosmosDBSourceReader(null, createConfig(37),
createRowType());
+
+ try {
+ Assertions.assertEquals(37, reader.getQueryPageSize());
+ } finally {
+ reader.close();
+ }
+ }
+
+ @Test
+ public void testSplitCopyPreservesContinuationToken() {
+ AzureCosmosDBSourceSplit split = new AzureCosmosDBSourceSplit(0,
"token-1");
+
+ AzureCosmosDBSourceSplit copiedSplit = split.copy();
+
+ Assertions.assertEquals(split.getSplitId(), copiedSplit.getSplitId());
+ Assertions.assertEquals(split.getContinuationToken(),
copiedSplit.getContinuationToken());
+ }
+
+ @Test
+ public void testRemoteFetchDoesNotHoldCheckpointLock() throws Exception {
+ CountDownLatch fetchStarted = new CountDownLatch(1);
+ CountDownLatch releaseFetch = new CountDownLatch(1);
+ RecordingCollector collector = new RecordingCollector();
+ BlockingFetchReader reader =
+ new BlockingFetchReader(
+ createConfig(1), createRowType(), fetchStarted,
releaseFetch);
+ AtomicReference<Throwable> pollFailure = new AtomicReference<>();
+ Thread pollThread =
+ new Thread(
+ () -> {
+ try {
+ reader.pollNext(collector);
+ } catch (Throwable throwable) {
+ pollFailure.set(throwable);
+ }
+ });
+
+ reader.addSplits(Collections.singletonList(new
AzureCosmosDBSourceSplit(0)));
+ pollThread.start();
+
+ Assertions.assertTrue(fetchStarted.await(5, TimeUnit.SECONDS));
+
+ ExecutorService checkpointThread = Executors.newSingleThreadExecutor();
+ try {
+ Future<Boolean> checkpointLockAcquired =
+ checkpointThread.submit(
+ () -> {
+ synchronized (collector.getCheckpointLock()) {
+ return true;
+ }
+ });
+ Assertions.assertTrue(checkpointLockAcquired.get(1,
TimeUnit.SECONDS));
+ } finally {
+ releaseFetch.countDown();
+ pollThread.join(TimeUnit.SECONDS.toMillis(5));
+ checkpointThread.shutdownNow();
+ reader.close();
+ }
+
+ Assertions.assertFalse(pollThread.isAlive());
+ if (pollFailure.get() != null) {
+ Assertions.fail(pollFailure.get());
+ }
+ }
+
+ private AzureCosmosDBConfig createConfig(int maxItemCount) {
+ Map<String, Object> schema = new HashMap<>();
+ schema.put("fields", new HashMap<String, Object>());
+
+ Map<String, Object> options = new HashMap<>();
+ options.put("endpoint", "https://account.documents.azure.com:443/");
+ options.put("primary_key", "account-key");
+ options.put("database", "sales");
+ options.put("container", "orders");
+ options.put("query", "SELECT * FROM c");
+ options.put("max_item_count", maxItemCount);
+ options.put("schema", schema);
+ return new AzureCosmosDBConfig(ReadonlyConfig.fromMap(options));
+ }
+
+ private SeaTunnelRowType createRowType() {
+ return new SeaTunnelRowType(
+ new String[] {"id"}, new SeaTunnelDataType[]
{BasicType.STRING_TYPE});
+ }
+
+ private static class BlockingFetchReader extends AzureCosmosDBSourceReader
{
+ private final CountDownLatch fetchStarted;
+ private final CountDownLatch releaseFetch;
+
+ private BlockingFetchReader(
+ AzureCosmosDBConfig config,
+ SeaTunnelRowType rowType,
+ CountDownLatch fetchStarted,
+ CountDownLatch releaseFetch) {
+ super(new RecordingReaderContext(), config, rowType);
+ this.fetchStarted = fetchStarted;
+ this.releaseFetch = releaseFetch;
+ }
+
+ @Override
+ FeedResponse<Object> fetchPage(String continuationToken) {
+ fetchStarted.countDown();
+ try {
+ if (!releaseFetch.await(5, TimeUnit.SECONDS)) {
+ throw new IllegalStateException("Timed out waiting to
release fetch");
+ }
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ throw new IllegalStateException("Interrupted while waiting to
release fetch", e);
+ }
+ return null;
+ }
+ }
+
+ private static class RecordingCollector implements Collector<SeaTunnelRow>
{
+ private final Object checkpointLock = new Object();
+
+ @Override
+ public void collect(SeaTunnelRow record) {}
+
+ @Override
+ public Object getCheckpointLock() {
+ return checkpointLock;
+ }
+ }
+
+ private static class RecordingReaderContext implements
SourceReader.Context {
+ @Override
+ public int getIndexOfSubtask() {
+ return 0;
+ }
+
+ @Override
+ public Boundedness getBoundedness() {
+ return Boundedness.BOUNDED;
+ }
+
+ @Override
+ public void signalNoMoreElement() {}
+
+ @Override
+ public void sendSplitRequest() {}
+
+ @Override
+ public void sendSourceEventToEnumerator(SourceEvent sourceEvent) {}
+
+ @Override
+ public MetricsContext getMetricsContext() {
+ return null;
+ }
+
+ @Override
+ public EventListener getEventListener() {
+ return null;
+ }
+ }
+}
diff --git
a/seatunnel-connectors-v2/connector-azurecosmosdb/src/test/java/org/apache/seatunnel/connectors/seatunnel/azurecosmosdb/source/AzureCosmosDBSourceSplitEnumeratorTest.java
b/seatunnel-connectors-v2/connector-azurecosmosdb/src/test/java/org/apache/seatunnel/connectors/seatunnel/azurecosmosdb/source/AzureCosmosDBSourceSplitEnumeratorTest.java
new file mode 100644
index 0000000000..ce85b07649
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-azurecosmosdb/src/test/java/org/apache/seatunnel/connectors/seatunnel/azurecosmosdb/source/AzureCosmosDBSourceSplitEnumeratorTest.java
@@ -0,0 +1,121 @@
+/*
+ * 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 org.apache.seatunnel.connectors.seatunnel.azurecosmosdb.source;
+
+import org.apache.seatunnel.api.common.metrics.MetricsContext;
+import org.apache.seatunnel.api.event.EventListener;
+import org.apache.seatunnel.api.source.SourceEvent;
+import org.apache.seatunnel.api.source.SourceSplitEnumerator;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.HashSet;
+import java.util.List;
+import java.util.Map;
+import java.util.Set;
+
+public class AzureCosmosDBSourceSplitEnumeratorTest {
+
+ @Test
+ public void testEnumeratorAssignsSingleSplit() throws Exception {
+ Set<Integer> readers = new HashSet<>(Arrays.asList(0, 1));
+ RecordingContext context = new RecordingContext(2, readers);
+ AzureCosmosDBSourceSplitEnumerator enumerator =
+ new AzureCosmosDBSourceSplitEnumerator(context, null);
+
+ try {
+ enumerator.run();
+ } finally {
+ enumerator.close();
+ }
+
+ Assertions.assertEquals(1, context.assignedSplits.get(0).size());
+ Assertions.assertEquals("0",
context.assignedSplits.get(0).get(0).splitId());
+ Assertions.assertTrue(
+ context.assignedSplits.getOrDefault(1,
Collections.emptyList()).isEmpty());
+ Assertions.assertEquals(readers, context.noMoreSplitReaders);
+ }
+
+ @Test
+ public void testAddSplitsBackPreservesContinuationToken() throws Exception
{
+ RecordingContext context = new RecordingContext(1,
Collections.singleton(0));
+ AzureCosmosDBSourceSplitEnumerator enumerator =
+ new AzureCosmosDBSourceSplitEnumerator(context, null);
+ AzureCosmosDBSourceSplit restoredSplit = new
AzureCosmosDBSourceSplit(0, "token-1");
+
+ try {
+ enumerator.addSplitsBack(Collections.singletonList(restoredSplit),
0);
+ } finally {
+ enumerator.close();
+ }
+
+ Assertions.assertEquals(
+ "token-1",
context.assignedSplits.get(0).get(0).getContinuationToken());
+ }
+
+ private static class RecordingContext
+ implements SourceSplitEnumerator.Context<AzureCosmosDBSourceSplit>
{
+
+ private final int parallelism;
+ private final Set<Integer> registeredReaders;
+ private final Map<Integer, List<AzureCosmosDBSourceSplit>>
assignedSplits = new HashMap<>();
+ private final Set<Integer> noMoreSplitReaders = new HashSet<>();
+
+ private RecordingContext(int parallelism, Set<Integer>
registeredReaders) {
+ this.parallelism = parallelism;
+ this.registeredReaders = registeredReaders;
+ }
+
+ @Override
+ public int currentParallelism() {
+ return parallelism;
+ }
+
+ @Override
+ public Set<Integer> registeredReaders() {
+ return registeredReaders;
+ }
+
+ @Override
+ public void assignSplit(int subtaskId, List<AzureCosmosDBSourceSplit>
splits) {
+ assignedSplits.put(subtaskId, splits);
+ }
+
+ @Override
+ public void signalNoMoreSplits(int subtask) {
+ noMoreSplitReaders.add(subtask);
+ }
+
+ @Override
+ public void sendEventToSourceReader(int subtaskId, SourceEvent event)
{}
+
+ @Override
+ public MetricsContext getMetricsContext() {
+ return null;
+ }
+
+ @Override
+ public EventListener getEventListener() {
+ return null;
+ }
+ }
+}
diff --git a/seatunnel-connectors-v2/pom.xml b/seatunnel-connectors-v2/pom.xml
index 694c90ad6c..c764541bb1 100644
--- a/seatunnel-connectors-v2/pom.xml
+++ b/seatunnel-connectors-v2/pom.xml
@@ -102,6 +102,7 @@
<module>connector-lance</module>
<module>connector-bigquery</module>
<module>connector-edge-socket</module>
+ <module>connector-azurecosmosdb</module>
<module>connector-nats-jetstream</module>
</modules>
diff --git a/seatunnel-dist/pom.xml b/seatunnel-dist/pom.xml
index 7354811388..aad8d7acc3 100644
--- a/seatunnel-dist/pom.xml
+++ b/seatunnel-dist/pom.xml
@@ -584,6 +584,12 @@
<version>${project.version}</version>
<scope>provided</scope>
</dependency>
+ <dependency>
+ <groupId>org.apache.seatunnel</groupId>
+ <artifactId>connector-azurecosmosdb</artifactId>
+ <version>${project.version}</version>
+ <scope>provided</scope>
+ </dependency>
<dependency>
<groupId>org.apache.seatunnel</groupId>
<artifactId>connector-starrocks</artifactId>
diff --git
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-azurecosmosdb-e2e/pom.xml
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-azurecosmosdb-e2e/pom.xml
new file mode 100644
index 0000000000..78c35471e7
--- /dev/null
+++
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-azurecosmosdb-e2e/pom.xml
@@ -0,0 +1,63 @@
+<?xml version="1.0" encoding="UTF-8"?>
+<!--
+ 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.
+-->
+<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
+ xsi:schemaLocation="http://maven.apache.org/POM/4.0.0
http://maven.apache.org/xsd/maven-4.0.0.xsd">
+ <modelVersion>4.0.0</modelVersion>
+ <parent>
+ <groupId>org.apache.seatunnel</groupId>
+ <artifactId>seatunnel-connector-v2-e2e</artifactId>
+ <version>${revision}</version>
+ </parent>
+
+ <artifactId>connector-azurecosmosdb-e2e</artifactId>
+ <name>SeaTunnel : E2E : Connector V2 : Azure Cosmos DB</name>
+
+ <dependencies>
+ <dependency>
+ <groupId>org.apache.seatunnel</groupId>
+ <artifactId>connector-azurecosmosdb</artifactId>
+ <version>${project.version}</version>
+ <scope>test</scope>
+ </dependency>
+ <dependency>
+ <groupId>org.apache.seatunnel</groupId>
+ <artifactId>connector-assert</artifactId>
+ <version>${project.version}</version>
+ <scope>test</scope>
+ </dependency>
+ <dependency>
+ <groupId>org.apache.seatunnel</groupId>
+ <artifactId>connector-common</artifactId>
+ <version>${project.version}</version>
+ <type>test-jar</type>
+ <scope>test</scope>
+ </dependency>
+ <dependency>
+ <groupId>com.azure</groupId>
+ <artifactId>azure-cosmos</artifactId>
+ <version>4.63.0</version>
+ <scope>test</scope>
+ </dependency>
+ <dependency>
+ <groupId>org.testcontainers</groupId>
+ <artifactId>azure</artifactId>
+ <version>${testcontainer.version}</version>
+ <scope>test</scope>
+ </dependency>
+ </dependencies>
+</project>
diff --git
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-azurecosmosdb-e2e/src/test/java/org/apache/seatunnel/e2e/connector/azurecosmosdb/AbstractAzureCosmosDBIT.java
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-azurecosmosdb-e2e/src/test/java/org/apache/seatunnel/e2e/connector/azurecosmosdb/AbstractAzureCosmosDBIT.java
new file mode 100644
index 0000000000..0c89ff8149
--- /dev/null
+++
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-azurecosmosdb-e2e/src/test/java/org/apache/seatunnel/e2e/connector/azurecosmosdb/AbstractAzureCosmosDBIT.java
@@ -0,0 +1,247 @@
+/*
+ * 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 org.apache.seatunnel.e2e.connector.azurecosmosdb;
+
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.source.Boundedness;
+import org.apache.seatunnel.api.source.Collector;
+import org.apache.seatunnel.api.source.SourceEvent;
+import org.apache.seatunnel.api.source.SourceReader;
+import org.apache.seatunnel.api.table.type.BasicType;
+import org.apache.seatunnel.api.table.type.SeaTunnelDataType;
+import org.apache.seatunnel.api.table.type.SeaTunnelRow;
+import org.apache.seatunnel.api.table.type.SeaTunnelRowType;
+import
org.apache.seatunnel.connectors.seatunnel.azurecosmosdb.config.AzureCosmosDBConfig;
+import
org.apache.seatunnel.connectors.seatunnel.azurecosmosdb.source.AzureCosmosDBSourceReader;
+import
org.apache.seatunnel.connectors.seatunnel.azurecosmosdb.source.AzureCosmosDBSourceSplit;
+import org.apache.seatunnel.e2e.common.TestResource;
+import org.apache.seatunnel.e2e.common.TestSuiteBase;
+
+import org.junit.jupiter.api.AfterAll;
+import org.junit.jupiter.api.BeforeAll;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+import org.testcontainers.containers.CosmosDBEmulatorContainer;
+import org.testcontainers.utility.DockerImageName;
+
+import com.azure.cosmos.CosmosClient;
+import com.azure.cosmos.CosmosClientBuilder;
+import com.azure.cosmos.CosmosContainer;
+import com.azure.cosmos.CosmosDatabase;
+import com.azure.cosmos.models.CosmosContainerProperties;
+
+import java.io.OutputStream;
+import java.nio.file.Files;
+import java.nio.file.Path;
+import java.security.KeyStore;
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+
+public abstract class AbstractAzureCosmosDBIT extends TestSuiteBase implements
TestResource {
+
+ private static final Logger LOG =
LoggerFactory.getLogger(AbstractAzureCosmosDBIT.class);
+
+ protected static final String DATABASE = "seatunnel_e2e";
+ protected static final String BASIC_CONTAINER = "source_basic_orders";
+ protected static final String FILTER_CONTAINER = "source_filter_orders";
+ protected static final String PAGINATION_CONTAINER =
"source_pagination_orders";
+ protected static final String TRUST_STORE_PASSWORD = "changeit";
+
+ private static final DockerImageName COSMOS_IMAGE =
+
DockerImageName.parse("mcr.microsoft.com/cosmosdb/linux/azure-cosmos-emulator:latest");
+
+ protected CosmosDBEmulatorContainer cosmosContainer;
+ protected CosmosClient client;
+
+ @BeforeAll
+ @Override
+ public void startUp() throws Exception {
+ cosmosContainer = new CosmosDBEmulatorContainer(COSMOS_IMAGE);
+ cosmosContainer.start();
+ installEmulatorTrustStore();
+ client =
+ new CosmosClientBuilder()
+ .endpoint(cosmosContainer.getEmulatorEndpoint())
+ .key(cosmosContainer.getEmulatorKey())
+ .endpointDiscoveryEnabled(false)
+ .gatewayMode()
+ .buildClient();
+ seedContainer(BASIC_CONTAINER, item("1", "alpha", 10), item("2",
"beta", 20));
+ seedContainer(
+ FILTER_CONTAINER,
+ item("1", "low-score", 5),
+ item("2", "high-score", 30),
+ item("3", "higher-score", 40));
+ seedContainer(
+ PAGINATION_CONTAINER,
+ item("1", "page-one", 10),
+ item("2", "page-two", 20),
+ item("3", "page-three", 30));
+ LOG.info("Azure Cosmos DB emulator started for source e2e tests");
+ }
+
+ @AfterAll
+ @Override
+ public void tearDown() {
+ if (client != null) {
+ client.close();
+ }
+ if (cosmosContainer != null) {
+ cosmosContainer.close();
+ }
+ }
+
+ protected List<SeaTunnelRow> readRows(String container, String query, int
maxItemCount)
+ throws Exception {
+ RecordingCollector collector = new RecordingCollector();
+ RecordingReaderContext context = new RecordingReaderContext();
+ AzureCosmosDBSourceReader reader =
+ new AzureCosmosDBSourceReader(
+ context, createConfig(container, query, maxItemCount),
rowType());
+
+ reader.open();
+ try {
+ reader.addSplits(Collections.singletonList(new
AzureCosmosDBSourceSplit(0)));
+ reader.handleNoMoreSplits();
+ while (!context.isNoMoreElement()) {
+ reader.pollNext(collector);
+ }
+ } finally {
+ reader.close();
+ }
+ return collector.rows;
+ }
+
+ protected AzureCosmosDBConfig createConfig(String container, String query,
int maxItemCount) {
+ Map<String, Object> schema = new HashMap<>();
+ schema.put("fields", new HashMap<String, Object>());
+
+ Map<String, Object> options = new HashMap<>();
+ options.put("endpoint", cosmosContainer.getEmulatorEndpoint());
+ options.put("key", cosmosContainer.getEmulatorKey());
+ options.put("database", DATABASE);
+ options.put("container", container);
+ options.put("query", query);
+ options.put("max_item_count", maxItemCount);
+ options.put("schema", schema);
+ return new AzureCosmosDBConfig(ReadonlyConfig.fromMap(options));
+ }
+
+ protected SeaTunnelRowType rowType() {
+ return new SeaTunnelRowType(
+ new String[] {"id", "name", "score"},
+ new SeaTunnelDataType[] {
+ BasicType.STRING_TYPE, BasicType.STRING_TYPE,
BasicType.INT_TYPE
+ });
+ }
+
+ protected static class RecordingCollector implements
Collector<SeaTunnelRow> {
+ private final List<SeaTunnelRow> rows = new ArrayList<>();
+ private final Object checkpointLock = new Object();
+
+ @Override
+ public void collect(SeaTunnelRow record) {
+ rows.add(record);
+ }
+
+ @Override
+ public Object getCheckpointLock() {
+ return checkpointLock;
+ }
+
+ public List<SeaTunnelRow> getRows() {
+ return rows;
+ }
+ }
+
+ protected static class RecordingReaderContext implements
SourceReader.Context {
+ private volatile boolean noMoreElement;
+
+ @Override
+ public int getIndexOfSubtask() {
+ return 0;
+ }
+
+ @Override
+ public Boundedness getBoundedness() {
+ return Boundedness.BOUNDED;
+ }
+
+ @Override
+ public void signalNoMoreElement() {
+ noMoreElement = true;
+ }
+
+ public boolean isNoMoreElement() {
+ return noMoreElement;
+ }
+
+ @Override
+ public void sendSplitRequest() {}
+
+ @Override
+ public void sendSourceEventToEnumerator(SourceEvent sourceEvent) {}
+
+ @Override
+ public org.apache.seatunnel.api.common.metrics.MetricsContext
getMetricsContext() {
+ return null;
+ }
+
+ @Override
+ public org.apache.seatunnel.api.event.EventListener getEventListener()
{
+ return null;
+ }
+ }
+
+ private void installEmulatorTrustStore() throws Exception {
+ KeyStore keyStore = cosmosContainer.buildNewKeyStore();
+ Path trustStore = Files.createTempFile("cosmos-emulator", ".jks");
+ try (OutputStream outputStream = Files.newOutputStream(trustStore)) {
+ keyStore.store(outputStream, TRUST_STORE_PASSWORD.toCharArray());
+ }
+ System.setProperty("javax.net.ssl.trustStore",
trustStore.toAbsolutePath().toString());
+ System.setProperty("javax.net.ssl.trustStorePassword",
TRUST_STORE_PASSWORD);
+ }
+
+ @SafeVarargs
+ private final void seedContainer(String containerName, Map<String,
Object>... items) {
+ client.createDatabaseIfNotExists(DATABASE);
+ CosmosDatabase database = client.getDatabase(DATABASE);
+ try {
+ database.getContainer(containerName).delete();
+ } catch (Exception ignored) {
+ // Container may not exist during first setup.
+ }
+ database.createContainerIfNotExists(new
CosmosContainerProperties(containerName, "/id"));
+ CosmosContainer container = database.getContainer(containerName);
+ for (Map<String, Object> item : items) {
+ container.upsertItem(item);
+ }
+ }
+
+ private Map<String, Object> item(String id, String name, int score) {
+ Map<String, Object> item = new HashMap<>();
+ item.put("id", id);
+ item.put("name", name);
+ item.put("score", score);
+ return item;
+ }
+}
diff --git
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-azurecosmosdb-e2e/src/test/java/org/apache/seatunnel/e2e/connector/azurecosmosdb/AzureCosmosDBSourceIT.java
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-azurecosmosdb-e2e/src/test/java/org/apache/seatunnel/e2e/connector/azurecosmosdb/AzureCosmosDBSourceIT.java
new file mode 100644
index 0000000000..c8693ae61e
--- /dev/null
+++
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-azurecosmosdb-e2e/src/test/java/org/apache/seatunnel/e2e/connector/azurecosmosdb/AzureCosmosDBSourceIT.java
@@ -0,0 +1,138 @@
+/*
+ * 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 org.apache.seatunnel.e2e.connector.azurecosmosdb;
+
+import org.apache.seatunnel.api.table.type.SeaTunnelRow;
+import
org.apache.seatunnel.connectors.seatunnel.azurecosmosdb.source.AzureCosmosDBSourceReader;
+import
org.apache.seatunnel.connectors.seatunnel.azurecosmosdb.source.AzureCosmosDBSourceSplit;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+import java.util.Collections;
+import java.util.Comparator;
+import java.util.HashSet;
+import java.util.List;
+import java.util.Set;
+
+public class AzureCosmosDBSourceIT extends AbstractAzureCosmosDBIT {
+
+ @Test
+ public void testBasicSourceRead() throws Exception {
+ List<SeaTunnelRow> rows =
+ readRows(BASIC_CONTAINER, "SELECT c.id, c.name, c.score FROM
c", 100);
+ rows.sort(Comparator.comparing(row ->
String.valueOf(row.getField(0))));
+
+ Assertions.assertEquals(2, rows.size());
+ assertRow(rows.get(0), "1", "alpha", 10);
+ assertRow(rows.get(1), "2", "beta", 20);
+ }
+
+ @Test
+ public void testSourceQueryFilterAndProjection() throws Exception {
+ List<SeaTunnelRow> rows =
+ readRows(
+ FILTER_CONTAINER,
+ "SELECT c.id, c.name, c.score FROM c WHERE c.score >
10",
+ 100);
+ rows.sort(Comparator.comparing(row ->
String.valueOf(row.getField(0))));
+
+ Assertions.assertEquals(2, rows.size());
+ assertRow(rows.get(0), "2", "high-score", 30);
+ assertRow(rows.get(1), "3", "higher-score", 40);
+ }
+
+ @Test
+ public void testSourcePaginationWithMaxItemCount() throws Exception {
+ List<SeaTunnelRow> rows =
+ readRows(PAGINATION_CONTAINER, "SELECT c.id, c.name, c.score
FROM c", 1);
+ rows.sort(Comparator.comparing(row ->
String.valueOf(row.getField(0))));
+
+ Assertions.assertEquals(3, rows.size());
+ assertRow(rows.get(0), "1", "page-one", 10);
+ assertRow(rows.get(1), "2", "page-two", 20);
+ assertRow(rows.get(2), "3", "page-three", 30);
+ }
+
+ @Test
+ public void testCheckpointRestoreResumesFromContinuationToken() throws
Exception {
+ String query = "SELECT c.id, c.name, c.score FROM c";
+ RecordingCollector firstCollector = new RecordingCollector();
+ RecordingReaderContext firstContext = new RecordingReaderContext();
+ AzureCosmosDBSourceReader firstReader =
+ new AzureCosmosDBSourceReader(
+ firstContext, createConfig(PAGINATION_CONTAINER,
query, 1), rowType());
+ List<AzureCosmosDBSourceSplit> checkpoint;
+
+ firstReader.open();
+ try {
+ firstReader.addSplits(Collections.singletonList(new
AzureCosmosDBSourceSplit(0)));
+ firstReader.handleNoMoreSplits();
+ firstReader.pollNext(firstCollector);
+ checkpoint = firstReader.snapshotState(1L);
+ } finally {
+ firstReader.close();
+ }
+
+ Assertions.assertEquals(1, firstCollector.getRows().size());
+ Assertions.assertEquals(1, checkpoint.size());
+ Assertions.assertNotNull(checkpoint.get(0).getContinuationToken());
+
+ RecordingCollector restoredCollector = new RecordingCollector();
+ RecordingReaderContext restoredContext = new RecordingReaderContext();
+ AzureCosmosDBSourceReader restoredReader =
+ new AzureCosmosDBSourceReader(
+ restoredContext, createConfig(PAGINATION_CONTAINER,
query, 1), rowType());
+
+ restoredReader.open();
+ try {
+ restoredReader.addSplits(checkpoint);
+ restoredReader.handleNoMoreSplits();
+ while (!restoredContext.isNoMoreElement()) {
+ restoredReader.pollNext(restoredCollector);
+ }
+ } finally {
+ restoredReader.close();
+ }
+
+ Set<String> firstReadIds = rowIds(firstCollector.getRows());
+ Set<String> restoredReadIds = rowIds(restoredCollector.getRows());
+ Set<String> allReadIds = new HashSet<>(firstReadIds);
+ allReadIds.addAll(restoredReadIds);
+
+ Assertions.assertEquals(2, restoredCollector.getRows().size());
+ for (String firstReadId : firstReadIds) {
+ Assertions.assertFalse(restoredReadIds.contains(firstReadId));
+ }
+ Assertions.assertEquals(3, allReadIds.size());
+ }
+
+ private void assertRow(SeaTunnelRow row, String id, String name, int
score) {
+ Assertions.assertEquals(id, row.getField(0));
+ Assertions.assertEquals(name, row.getField(1));
+ Assertions.assertEquals(score, row.getField(2));
+ }
+
+ private Set<String> rowIds(List<SeaTunnelRow> rows) {
+ Set<String> ids = new HashSet<>();
+ for (SeaTunnelRow row : rows) {
+ ids.add(String.valueOf(row.getField(0)));
+ }
+ return ids;
+ }
+}
diff --git
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-azurecosmosdb-e2e/src/test/resources/azurecosmosdb/azurecosmosdb_source_basic.conf
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-azurecosmosdb-e2e/src/test/resources/azurecosmosdb/azurecosmosdb_source_basic.conf
new file mode 100644
index 0000000000..71a5bd16c4
--- /dev/null
+++
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-azurecosmosdb-e2e/src/test/resources/azurecosmosdb/azurecosmosdb_source_basic.conf
@@ -0,0 +1,73 @@
+#
+# 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.
+#
+
+env {
+ parallelism = 1
+ job.mode = "BATCH"
+}
+
+source {
+ AzureCosmosDB {
+ endpoint = "${cosmos_endpoint}"
+ key = "${cosmos_key}"
+ database = "seatunnel_e2e"
+ container = "source_basic_orders"
+ query = "SELECT c.id, c.name, c.score FROM c"
+ max_item_count = 100
+ schema = {
+ fields {
+ id = string
+ name = string
+ score = int
+ }
+ }
+ }
+}
+
+sink {
+ Assert {
+ rules {
+ row_rules = [
+ {
+ rule_type = MIN_ROW
+ rule_value = 2
+ },
+ {
+ rule_type = MAX_ROW
+ rule_value = 2
+ }
+ ],
+ field_rules = [
+ {
+ field_name = id
+ field_type = string
+ field_value = [{rule_type = NOT_NULL}]
+ },
+ {
+ field_name = name
+ field_type = string
+ field_value = [{rule_type = NOT_NULL}]
+ },
+ {
+ field_name = score
+ field_type = int
+ field_value = [{rule_type = MIN, rule_value = 10}]
+ }
+ ]
+ }
+ }
+}
diff --git
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-azurecosmosdb-e2e/src/test/resources/azurecosmosdb/azurecosmosdb_source_pagination.conf
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-azurecosmosdb-e2e/src/test/resources/azurecosmosdb/azurecosmosdb_source_pagination.conf
new file mode 100644
index 0000000000..3675e73ab9
--- /dev/null
+++
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-azurecosmosdb-e2e/src/test/resources/azurecosmosdb/azurecosmosdb_source_pagination.conf
@@ -0,0 +1,68 @@
+#
+# 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.
+#
+
+env {
+ parallelism = 1
+ job.mode = "BATCH"
+}
+
+source {
+ AzureCosmosDB {
+ endpoint = "${cosmos_endpoint}"
+ key = "${cosmos_key}"
+ database = "seatunnel_e2e"
+ container = "source_pagination_orders"
+ query = "SELECT c.id, c.name, c.score FROM c"
+ max_item_count = 1
+ schema = {
+ fields {
+ id = string
+ name = string
+ score = int
+ }
+ }
+ }
+}
+
+sink {
+ Assert {
+ rules {
+ row_rules = [
+ {
+ rule_type = MIN_ROW
+ rule_value = 3
+ },
+ {
+ rule_type = MAX_ROW
+ rule_value = 3
+ }
+ ],
+ field_rules = [
+ {
+ field_name = id
+ field_type = string
+ field_value = [{rule_type = NOT_NULL}]
+ },
+ {
+ field_name = score
+ field_type = int
+ field_value = [{rule_type = MIN, rule_value = 10}]
+ }
+ ]
+ }
+ }
+}
diff --git
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-azurecosmosdb-e2e/src/test/resources/azurecosmosdb/azurecosmosdb_source_query_filter.conf
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-azurecosmosdb-e2e/src/test/resources/azurecosmosdb/azurecosmosdb_source_query_filter.conf
new file mode 100644
index 0000000000..abacf3959a
--- /dev/null
+++
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-azurecosmosdb-e2e/src/test/resources/azurecosmosdb/azurecosmosdb_source_query_filter.conf
@@ -0,0 +1,68 @@
+#
+# 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.
+#
+
+env {
+ parallelism = 1
+ job.mode = "BATCH"
+}
+
+source {
+ AzureCosmosDB {
+ endpoint = "${cosmos_endpoint}"
+ key = "${cosmos_key}"
+ database = "seatunnel_e2e"
+ container = "source_filter_orders"
+ query = "SELECT c.id, c.name, c.score FROM c WHERE c.score > 10"
+ max_item_count = 100
+ schema = {
+ fields {
+ id = string
+ name = string
+ score = int
+ }
+ }
+ }
+}
+
+sink {
+ Assert {
+ rules {
+ row_rules = [
+ {
+ rule_type = MIN_ROW
+ rule_value = 2
+ },
+ {
+ rule_type = MAX_ROW
+ rule_value = 2
+ }
+ ],
+ field_rules = [
+ {
+ field_name = score
+ field_type = int
+ field_value = [{rule_type = MIN, rule_value = 11}]
+ },
+ {
+ field_name = name
+ field_type = string
+ field_value = [{rule_type = NOT_NULL}]
+ }
+ ]
+ }
+ }
+}
diff --git a/seatunnel-e2e/seatunnel-connector-v2-e2e/pom.xml
b/seatunnel-e2e/seatunnel-connector-v2-e2e/pom.xml
index 27e2487aa2..61324e27d3 100644
--- a/seatunnel-e2e/seatunnel-connector-v2-e2e/pom.xml
+++ b/seatunnel-e2e/seatunnel-connector-v2-e2e/pom.xml
@@ -37,6 +37,7 @@
<module>connector-starrocks-e2e</module>
<module>connector-influxdb-e2e</module>
<module>connector-amazondynamodb-e2e</module>
+ <module>connector-azurecosmosdb-e2e</module>
<module>connector-amazonsqs-e2e</module>
<module>connector-file-local-e2e</module>
<module>connector-file-cos-e2e</module>