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>

Reply via email to