This is an automated email from the ASF dual-hosted git repository.

github-merge-queue[bot] pushed a commit to branch 
gh-readonly-queue/dev/pr-11998-c5b58851ba9e5aed7ea5d7cee872952616df814c
in repository https://gitbox.apache.org/repos/asf/seatunnel.git

commit e352beba2911cbd562a73e2627f0153a99fd0209
Author: Goutam Adwant <[email protected]>
AuthorDate: Sun Aug 30 15:56:34 2026 +0000

    [Feature][Connector-V2] Add PostHog source connector (#11998)
    
    Signed-off-by: goutamadwant <[email protected]>
---
 config/plugin_config                               |   1 +
 .../connectors/changelog/connector-http-posthog.md |   6 +
 docs/en/connectors/source/PostHog.md               |  76 ++++++
 .../connectors/changelog/connector-http-posthog.md |   6 +
 docs/zh/connectors/source/PostHog.md               |  76 ++++++
 plugin-mapping.properties                          |   1 +
 .../{ => connector-http-posthog}/pom.xml           |  33 +--
 .../seatunnel/posthog/source/PostHogSource.java    |  82 ++++++
 .../posthog/source/PostHogSourceFactory.java       |  73 +++++
 .../posthog/source/PostHogSourceReader.java        | 206 +++++++++++++++
 .../source/config/PostHogSourceOptions.java        |  50 ++++
 .../source/config/PostHogSourceParameter.java      | 110 ++++++++
 .../posthog/source/PostHogSourceFactoryTest.java   |  33 +++
 .../posthog/source/PostHogSourceReaderTest.java    | 294 +++++++++++++++++++++
 .../source/config/PostHogSourceParameterTest.java  | 105 ++++++++
 seatunnel-connectors-v2/connector-http/pom.xml     |   1 +
 seatunnel-dist/pom.xml                             |   6 +
 .../connector-http-e2e/pom.xml                     |   6 +
 .../seatunnel/e2e/connector/http/HttpIT.java       |   4 +
 .../src/test/resources/mockserver-config.json      |  30 ++-
 .../src/test/resources/posthog_json_to_assert.conf |  67 +++++
 21 files changed, 1243 insertions(+), 23 deletions(-)

diff --git a/config/plugin_config b/config/plugin_config
index 06474d2733..99b7385da6 100644
--- a/config/plugin_config
+++ b/config/plugin_config
@@ -67,6 +67,7 @@ connector-http-notion
 connector-http-onesignal
 connector-http-wechat
 connector-http-airtable
+connector-http-posthog
 connector-hudi
 connector-iceberg
 connector-influxdb
diff --git a/docs/en/connectors/changelog/connector-http-posthog.md 
b/docs/en/connectors/changelog/connector-http-posthog.md
new file mode 100644
index 0000000000..1cf41e674b
--- /dev/null
+++ b/docs/en/connectors/changelog/connector-http-posthog.md
@@ -0,0 +1,6 @@
+<details><summary> Change Log </summary>
+
+| Change | Commit | Version |
+| --- | --- | --- |
+
+</details>
diff --git a/docs/en/connectors/source/PostHog.md 
b/docs/en/connectors/source/PostHog.md
new file mode 100644
index 0000000000..7f238d10df
--- /dev/null
+++ b/docs/en/connectors/source/PostHog.md
@@ -0,0 +1,76 @@
+import ChangeLog from '../changelog/connector-http-posthog.md';
+
+# PostHog
+
+> PostHog source connector
+
+## Description
+
+Reads the result of one HogQL query from PostHog. The connector uses the 
synchronous Query API and runs as a bounded batch source.
+
+For large historical exports, use [PostHog Batch 
Exports](https://posthog.com/docs/cdp/batch-exports) instead of issuing one 
large HogQL query.
+
+## 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)
+- [x] [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)
+
+## Source Options
+
+| Name | Type | Required | Default | Description |
+| --- | --- | --- | --- | --- |
+| base_url | String | No | https://us.posthog.com | PostHog instance base URL. 
Use `https://eu.posthog.com` for PostHog EU Cloud, or the URL of a self-hosted 
instance. |
+| project_id | String | Yes | - | PostHog project ID. |
+| api_key | String | Yes | - | PostHog personal API key with `query:read` 
permission. |
+| query | String | Yes | - | HogQL query executed through the PostHog Query 
API. |
+| schema | Config | Yes | - | Output schema. Every schema field must match a 
returned HogQL column name or alias. |
+| headers | Map | No | - | Additional HTTP headers. The connector sets the 
authorization and JSON headers. |
+| retry | int | No | 0 | Maximum number of HTTP request attempts after I/O 
failures. |
+| retry_backoff_multiplier_ms | int | No | 100 | Retry backoff multiplier in 
milliseconds. |
+| retry_backoff_max_ms | int | No | 10000 | Maximum retry backoff in 
milliseconds. |
+| connect_timeout_ms | int | No | 12000 | HTTP connection timeout in 
milliseconds. |
+| socket_timeout_ms | int | No | 60000 | HTTP socket timeout in milliseconds. |
+| common-options | Config | No | - | Source plugin common parameters. See 
[Source Common Options](../common-options/source-common-options.md). |
+
+## Usage Notes
+
+- The connector executes the configured query once and then finishes. Add a 
`LIMIT` and appropriate time filters so that the result fits in one PostHog 
Query API response.
+- HogQL expressions may return generated column names. Use `AS` aliases so 
every selected column matches a field in `schema`.
+- Keep `api_key` outside shared configuration files by using SeaTunnel 
variable substitution or the secret mechanism of the deployment platform.
+- The connector does not use PostHog's deprecated events-list endpoint and 
does not rewrite the supplied HogQL query.
+
+## Task Example
+
+```hocon
+env {
+  parallelism = 1
+  job.mode = "BATCH"
+}
+
+source {
+  PostHog {
+    base_url = "https://us.posthog.com";
+    project_id = "12345"
+    api_key = "${POSTHOG_API_KEY}"
+    query = "SELECT event, distinct_id, timestamp FROM events WHERE timestamp 
>= now() - INTERVAL 1 DAY ORDER BY timestamp LIMIT 10000"
+    schema = {
+      fields {
+        event = string
+        distinct_id = string
+        timestamp = timestamp
+      }
+    }
+  }
+}
+
+sink {
+  Console {
+  }
+}
+```
+
+<ChangeLog />
diff --git a/docs/zh/connectors/changelog/connector-http-posthog.md 
b/docs/zh/connectors/changelog/connector-http-posthog.md
new file mode 100644
index 0000000000..1cf41e674b
--- /dev/null
+++ b/docs/zh/connectors/changelog/connector-http-posthog.md
@@ -0,0 +1,6 @@
+<details><summary> Change Log </summary>
+
+| Change | Commit | Version |
+| --- | --- | --- |
+
+</details>
diff --git a/docs/zh/connectors/source/PostHog.md 
b/docs/zh/connectors/source/PostHog.md
new file mode 100644
index 0000000000..70b824ad05
--- /dev/null
+++ b/docs/zh/connectors/source/PostHog.md
@@ -0,0 +1,76 @@
+import ChangeLog from '../changelog/connector-http-posthog.md';
+
+# PostHog
+
+> PostHog 源连接器
+
+## 描述
+
+从 PostHog 读取一次 HogQL 查询结果。该连接器使用同步 Query API,并作为有界批处理源运行。
+
+对于大规模历史数据导出,请使用 [PostHog Batch 
Exports](https://posthog.com/docs/cdp/batch-exports),不要执行单个超大 HogQL 查询。
+
+## 主要特性
+
+- [x] [批](../../introduction/concepts/connector-v2-features.md)
+- [ ] [流](../../introduction/concepts/connector-v2-features.md)
+- [ ] [精确一次](../../introduction/concepts/connector-v2-features.md)
+- [x] [列投影](../../introduction/concepts/connector-v2-features.md)
+- [ ] [并行度](../../introduction/concepts/connector-v2-features.md)
+- [ ] [支持用户自定义分片](../../introduction/concepts/connector-v2-features.md)
+
+## 源选项
+
+| 参数名 | 类型 | 必须 | 默认值 | 描述 |
+| --- | --- | --- | --- | --- |
+| base_url | String | 否 | https://us.posthog.com | PostHog 实例基础 URL。PostHog 
欧洲云使用 `https://eu.posthog.com`,自托管部署使用对应实例 URL。 |
+| project_id | String | 是 | - | PostHog 项目 ID。 |
+| api_key | String | 是 | - | 具有 `query:read` 权限的 PostHog 个人 API 密钥。 |
+| query | String | 是 | - | 通过 PostHog Query API 执行的 HogQL 查询。 |
+| schema | Config | 是 | - | 输出结构。每个字段名必须与 HogQL 返回的列名或别名匹配。 |
+| headers | Map | 否 | - | 额外的 HTTP 请求头。连接器会设置认证头和 JSON 请求头。 |
+| retry | int | 否 | 0 | I/O 失败后的最大 HTTP 请求尝试次数。 |
+| retry_backoff_multiplier_ms | int | 否 | 100 | 重试退避乘数,单位毫秒。 |
+| retry_backoff_max_ms | int | 否 | 10000 | 最大重试退避时间,单位毫秒。 |
+| connect_timeout_ms | int | 否 | 12000 | HTTP 连接超时时间,单位毫秒。 |
+| socket_timeout_ms | int | 否 | 60000 | HTTP 套接字超时时间,单位毫秒。 |
+| common-options | Config | 否 | - | 
源插件通用参数,详见[源通用选项](../common-options/source-common-options.md)。 |
+
+## 使用提示
+
+- 连接器只执行一次配置的查询,然后结束。请添加 `LIMIT` 和适当的时间过滤条件,确保结果可由一次 PostHog Query API 响应返回。
+- HogQL 表达式可能生成默认列名。请使用 `AS` 别名,使每个返回列与 `schema` 字段匹配。
+- 请通过 SeaTunnel 变量替换或部署平台的密钥管理机制提供 `api_key`,不要在共享配置中保存真实密钥。
+- 连接器不使用 PostHog 已弃用的事件列表接口,也不会改写用户提供的 HogQL 查询。
+
+## 任务示例
+
+```hocon
+env {
+  parallelism = 1
+  job.mode = "BATCH"
+}
+
+source {
+  PostHog {
+    base_url = "https://us.posthog.com";
+    project_id = "12345"
+    api_key = "${POSTHOG_API_KEY}"
+    query = "SELECT event, distinct_id, timestamp FROM events WHERE timestamp 
>= now() - INTERVAL 1 DAY ORDER BY timestamp LIMIT 10000"
+    schema = {
+      fields {
+        event = string
+        distinct_id = string
+        timestamp = timestamp
+      }
+    }
+  }
+}
+
+sink {
+  Console {
+  }
+}
+```
+
+<ChangeLog />
diff --git a/plugin-mapping.properties b/plugin-mapping.properties
index 066de0c6ce..65a12260d2 100644
--- a/plugin-mapping.properties
+++ b/plugin-mapping.properties
@@ -109,6 +109,7 @@ seatunnel.source.Github = connector-http-github
 seatunnel.source.Notion = connector-http-notion
 seatunnel.source.Airtable = connector-http-airtable
 seatunnel.sink.Airtable = connector-http-airtable
+seatunnel.source.PostHog = connector-http-posthog
 seatunnel.sink.RabbitMQ = connector-rabbitmq
 seatunnel.source.RabbitMQ = connector-rabbitmq
 seatunnel.source.OpenMldb = connector-openmldb
diff --git a/seatunnel-connectors-v2/connector-http/pom.xml 
b/seatunnel-connectors-v2/connector-http/connector-http-posthog/pom.xml
similarity index 58%
copy from seatunnel-connectors-v2/connector-http/pom.xml
copy to seatunnel-connectors-v2/connector-http/connector-http-posthog/pom.xml
index a3bddc1c59..59dd8ecfad 100644
--- a/seatunnel-connectors-v2/connector-http/pom.xml
+++ b/seatunnel-connectors-v2/connector-http/connector-http-posthog/pom.xml
@@ -1,6 +1,5 @@
 <?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.
@@ -15,35 +14,25 @@
     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>
+        <artifactId>connector-http</artifactId>
         <version>${revision}</version>
     </parent>
-    <artifactId>connector-http</artifactId>
-    <packaging>pom</packaging>
-    <name>SeaTunnel : Connectors V2 : Http :</name>
 
-    <modules>
-        <module>connector-http-base</module>
-        <module>connector-http-feishu</module>
-        <module>connector-http-wechat</module>
-        <module>connector-http-myhours</module>
-        <module>connector-http-lemlist</module>
-        <module>connector-http-klaviyo</module>
-        <module>connector-http-onesignal</module>
-        <module>connector-http-jira</module>
-        <module>connector-http-gitlab</module>
-        <module>connector-http-github</module>
-        <module>connector-http-notion</module>
-        <module>connector-http-persistiq</module>
-        <module>connector-http-airtable</module>
-        <module>connector-http-zendesk</module>
-    </modules>
+    <artifactId>connector-http-posthog</artifactId>
+    <name>SeaTunnel : Connectors V2 : Http : PostHog</name>
+
+    <dependencies>
+        <dependency>
+            <groupId>org.apache.seatunnel</groupId>
+            <artifactId>connector-http-base</artifactId>
+            <version>${project.version}</version>
+        </dependency>
+    </dependencies>
 
 </project>
diff --git 
a/seatunnel-connectors-v2/connector-http/connector-http-posthog/src/main/java/org/apache/seatunnel/connectors/seatunnel/posthog/source/PostHogSource.java
 
b/seatunnel-connectors-v2/connector-http/connector-http-posthog/src/main/java/org/apache/seatunnel/connectors/seatunnel/posthog/source/PostHogSource.java
new file mode 100644
index 0000000000..4aea93298d
--- /dev/null
+++ 
b/seatunnel-connectors-v2/connector-http/connector-http-posthog/src/main/java/org/apache/seatunnel/connectors/seatunnel/posthog/source/PostHogSource.java
@@ -0,0 +1,82 @@
+/*
+ * 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.posthog.source;
+
+import org.apache.seatunnel.api.common.JobContext;
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.source.Boundedness;
+import org.apache.seatunnel.api.table.catalog.CatalogTable;
+import org.apache.seatunnel.api.table.catalog.CatalogTableUtil;
+import org.apache.seatunnel.api.table.type.SeaTunnelRow;
+import org.apache.seatunnel.common.constants.JobMode;
+import 
org.apache.seatunnel.connectors.seatunnel.common.source.AbstractSingleSplitReader;
+import 
org.apache.seatunnel.connectors.seatunnel.common.source.AbstractSingleSplitSource;
+import 
org.apache.seatunnel.connectors.seatunnel.common.source.SingleSplitReaderContext;
+import 
org.apache.seatunnel.connectors.seatunnel.posthog.source.config.PostHogSourceParameter;
+import org.apache.seatunnel.format.json.JsonDeserializationSchema;
+
+import java.util.Collections;
+import java.util.List;
+
+public class PostHogSource extends AbstractSingleSplitSource<SeaTunnelRow> {
+
+    public static final String PLUGIN_NAME = "PostHog";
+
+    private final PostHogSourceParameter sourceParameter = new 
PostHogSourceParameter();
+    private final CatalogTable catalogTable;
+    private JobContext jobContext;
+
+    public PostHogSource(ReadonlyConfig pluginConfig) {
+        sourceParameter.buildWithConfig(pluginConfig);
+        catalogTable = CatalogTableUtil.buildWithConfig(pluginConfig);
+    }
+
+    @Override
+    public String getPluginName() {
+        return PLUGIN_NAME;
+    }
+
+    @Override
+    public void setJobContext(JobContext jobContext) {
+        this.jobContext = jobContext;
+    }
+
+    @Override
+    public Boundedness getBoundedness() {
+        if (!JobMode.BATCH.equals(jobContext.getJobMode())) {
+            throw new UnsupportedOperationException(
+                    "PostHog source connector only supports batch mode");
+        }
+        return Boundedness.BOUNDED;
+    }
+
+    @Override
+    public List<CatalogTable> getProducedCatalogTables() {
+        return Collections.singletonList(catalogTable);
+    }
+
+    @Override
+    public AbstractSingleSplitReader<SeaTunnelRow> createReader(
+            SingleSplitReaderContext readerContext) {
+        return new PostHogSourceReader(
+                sourceParameter,
+                readerContext,
+                new JsonDeserializationSchema(catalogTable, false, false),
+                catalogTable.getSeaTunnelRowType());
+    }
+}
diff --git 
a/seatunnel-connectors-v2/connector-http/connector-http-posthog/src/main/java/org/apache/seatunnel/connectors/seatunnel/posthog/source/PostHogSourceFactory.java
 
b/seatunnel-connectors-v2/connector-http/connector-http-posthog/src/main/java/org/apache/seatunnel/connectors/seatunnel/posthog/source/PostHogSourceFactory.java
new file mode 100644
index 0000000000..495e1f198e
--- /dev/null
+++ 
b/seatunnel-connectors-v2/connector-http/connector-http-posthog/src/main/java/org/apache/seatunnel/connectors/seatunnel/posthog/source/PostHogSourceFactory.java
@@ -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.
+ */
+
+package org.apache.seatunnel.connectors.seatunnel.posthog.source;
+
+import org.apache.seatunnel.api.configuration.util.OptionRule;
+import org.apache.seatunnel.api.options.ConnectorCommonOptions;
+import org.apache.seatunnel.api.source.SeaTunnelSource;
+import org.apache.seatunnel.api.source.SourceSplit;
+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.http.config.HttpCommonOptions;
+import org.apache.seatunnel.connectors.seatunnel.http.config.HttpSourceOptions;
+import 
org.apache.seatunnel.connectors.seatunnel.posthog.source.config.PostHogSourceOptions;
+
+import com.google.auto.service.AutoService;
+
+import java.io.Serializable;
+
+@AutoService(Factory.class)
+public class PostHogSourceFactory implements TableSourceFactory {
+
+    @Override
+    public String factoryIdentifier() {
+        return PostHogSource.PLUGIN_NAME;
+    }
+
+    @Override
+    public OptionRule optionRule() {
+        return OptionRule.builder()
+                .required(
+                        PostHogSourceOptions.PROJECT_ID,
+                        PostHogSourceOptions.API_KEY,
+                        PostHogSourceOptions.QUERY,
+                        ConnectorCommonOptions.SCHEMA)
+                .optional(
+                        PostHogSourceOptions.BASE_URL,
+                        HttpCommonOptions.HEADERS,
+                        HttpCommonOptions.RETRY,
+                        HttpCommonOptions.RETRY_BACKOFF_MULTIPLIER_MS,
+                        HttpCommonOptions.RETRY_BACKOFF_MAX_MS,
+                        HttpSourceOptions.CONNECT_TIMEOUT_MS,
+                        HttpSourceOptions.SOCKET_TIMEOUT_MS)
+                .build();
+    }
+
+    @Override
+    public <T, SplitT extends SourceSplit, StateT extends Serializable>
+            TableSource<T, SplitT, StateT> 
createSource(TableSourceFactoryContext context) {
+        return () -> (SeaTunnelSource<T, SplitT, StateT>) new 
PostHogSource(context.getOptions());
+    }
+
+    @Override
+    public Class<? extends SeaTunnelSource> getSourceClass() {
+        return PostHogSource.class;
+    }
+}
diff --git 
a/seatunnel-connectors-v2/connector-http/connector-http-posthog/src/main/java/org/apache/seatunnel/connectors/seatunnel/posthog/source/PostHogSourceReader.java
 
b/seatunnel-connectors-v2/connector-http/connector-http-posthog/src/main/java/org/apache/seatunnel/connectors/seatunnel/posthog/source/PostHogSourceReader.java
new file mode 100644
index 0000000000..c775d72453
--- /dev/null
+++ 
b/seatunnel-connectors-v2/connector-http/connector-http-posthog/src/main/java/org/apache/seatunnel/connectors/seatunnel/posthog/source/PostHogSourceReader.java
@@ -0,0 +1,206 @@
+/*
+ * 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.posthog.source;
+
+import org.apache.seatunnel.shade.com.fasterxml.jackson.databind.JsonNode;
+import 
org.apache.seatunnel.shade.com.fasterxml.jackson.databind.node.ObjectNode;
+import 
org.apache.seatunnel.shade.com.google.common.annotations.VisibleForTesting;
+
+import org.apache.seatunnel.api.serialization.DeserializationSchema;
+import org.apache.seatunnel.api.source.Collector;
+import org.apache.seatunnel.api.table.type.SeaTunnelRow;
+import org.apache.seatunnel.api.table.type.SeaTunnelRowType;
+import org.apache.seatunnel.common.utils.JsonUtils;
+import 
org.apache.seatunnel.connectors.seatunnel.common.source.AbstractSingleSplitReader;
+import 
org.apache.seatunnel.connectors.seatunnel.common.source.SingleSplitReaderContext;
+import 
org.apache.seatunnel.connectors.seatunnel.http.client.HttpClientProvider;
+import org.apache.seatunnel.connectors.seatunnel.http.client.HttpResponse;
+import org.apache.seatunnel.connectors.seatunnel.http.config.HttpParameter;
+import 
org.apache.seatunnel.connectors.seatunnel.http.exception.HttpConnectorErrorCode;
+import 
org.apache.seatunnel.connectors.seatunnel.http.exception.HttpConnectorException;
+import 
org.apache.seatunnel.connectors.seatunnel.http.source.DeserializationCollector;
+
+import java.io.IOException;
+import java.nio.charset.StandardCharsets;
+import java.util.Collections;
+import java.util.HashSet;
+import java.util.LinkedHashSet;
+import java.util.Set;
+
+public class PostHogSourceReader extends 
AbstractSingleSplitReader<SeaTunnelRow> {
+
+    private final HttpParameter httpParameter;
+    private final SingleSplitReaderContext context;
+    private final DeserializationCollector deserializationCollector;
+    private final SeaTunnelRowType rowType;
+    private HttpClientProvider httpClient;
+
+    public PostHogSourceReader(
+            HttpParameter httpParameter,
+            SingleSplitReaderContext context,
+            DeserializationSchema<SeaTunnelRow> deserializationSchema,
+            SeaTunnelRowType rowType) {
+        this(httpParameter, context, new 
DeserializationCollector(deserializationSchema), rowType);
+    }
+
+    @VisibleForTesting
+    PostHogSourceReader(
+            HttpParameter httpParameter,
+            SingleSplitReaderContext context,
+            DeserializationCollector deserializationCollector,
+            SeaTunnelRowType rowType) {
+        this.httpParameter = httpParameter;
+        this.context = context;
+        this.deserializationCollector = deserializationCollector;
+        this.rowType = rowType;
+    }
+
+    @Override
+    public void open() {
+        httpClient = new HttpClientProvider(httpParameter);
+    }
+
+    @VisibleForTesting
+    void setHttpClient(HttpClientProvider httpClient) {
+        this.httpClient = httpClient;
+    }
+
+    @Override
+    public void close() throws IOException {
+        if (httpClient != null) {
+            httpClient.close();
+        }
+    }
+
+    @Override
+    public void internalPollNext(Collector<SeaTunnelRow> output) throws 
Exception {
+        HttpResponse response =
+                httpClient.doPost(
+                        httpParameter.getUrl(),
+                        httpParameter.getHeaders(),
+                        httpParameter.getBody());
+        if (response.getCode() < 200 || response.getCode() >= 300) {
+            throw requestFailed(
+                    "PostHog query request failed with HTTP status " + 
response.getCode());
+        }
+        collectResponse(response.getContent(), output);
+        context.signalNoMoreElement();
+    }
+
+    @VisibleForTesting
+    void collectResponse(String content, Collector<SeaTunnelRow> output) 
throws IOException {
+        if (content == null || content.trim().isEmpty()) {
+            throw requestFailed("PostHog query returned an empty response");
+        }
+
+        JsonNode response;
+        try {
+            response = JsonUtils.stringToJsonNode(content);
+        } catch (Exception e) {
+            throw new HttpConnectorException(
+                    HttpConnectorErrorCode.REQUEST_FAILED,
+                    "PostHog query returned invalid JSON",
+                    e);
+        }
+        validateQueryStatus(response);
+
+        JsonNode columnsNode = response.get("columns");
+        JsonNode resultsNode = response.get("results");
+        if (columnsNode == null || !columnsNode.isArray()) {
+            throw requestFailed("PostHog query response is missing the columns 
array");
+        }
+        if (resultsNode == null || !resultsNode.isArray()) {
+            throw requestFailed("PostHog query response is missing the results 
array");
+        }
+
+        String[] columns = readColumns(columnsNode);
+        validateSchemaColumns(columns);
+        for (JsonNode result : resultsNode) {
+            if (!result.isArray()) {
+                throw requestFailed("PostHog query result rows must be 
arrays");
+            }
+            if (result.size() != columns.length) {
+                throw requestFailed("PostHog query result width does not match 
the columns array");
+            }
+            ObjectNode row = JsonUtils.createObjectNode();
+            for (int index = 0; index < columns.length; index++) {
+                row.set(columns[index], result.get(index));
+            }
+            deserializationCollector.collect(
+                    row.toString().getBytes(StandardCharsets.UTF_8), output);
+        }
+    }
+
+    private void validateQueryStatus(JsonNode response) {
+        JsonNode error = response.get("error");
+        if (error != null && !error.isNull() && 
!error.asText().trim().isEmpty()) {
+            throw requestFailed("PostHog query failed: " + error.asText());
+        }
+
+        JsonNode queryStatus = response.get("query_status");
+        if (queryStatus == null || queryStatus.isNull()) {
+            return;
+        }
+        if (queryStatus.path("error").asBoolean(false)) {
+            String message = queryStatus.path("error_message").asText("unknown 
query error");
+            throw requestFailed("PostHog query failed: " + message);
+        }
+        if (queryStatus.has("complete") && 
!queryStatus.path("complete").asBoolean()) {
+            throw requestFailed("PostHog query did not complete in blocking 
mode");
+        }
+    }
+
+    private String[] readColumns(JsonNode columnsNode) {
+        String[] columns = new String[columnsNode.size()];
+        Set<String> uniqueColumns = new HashSet<>();
+        for (int index = 0; index < columnsNode.size(); index++) {
+            JsonNode columnNode = columnsNode.get(index);
+            if (!columnNode.isTextual() || 
columnNode.asText().trim().isEmpty()) {
+                throw requestFailed("PostHog query response contains an 
invalid column name");
+            }
+            String column = columnNode.asText();
+            if (!uniqueColumns.add(column)) {
+                throw requestFailed(
+                        "PostHog query returned duplicate column name '" + 
column + "'");
+            }
+            columns[index] = column;
+        }
+        return columns;
+    }
+
+    private void validateSchemaColumns(String[] columns) {
+        Set<String> availableColumns = new LinkedHashSet<>();
+        Collections.addAll(availableColumns, columns);
+        Set<String> missingColumns = new LinkedHashSet<>();
+        for (String fieldName : rowType.getFieldNames()) {
+            if (!availableColumns.contains(fieldName)) {
+                missingColumns.add(fieldName);
+            }
+        }
+        if (!missingColumns.isEmpty()) {
+            throw requestFailed(
+                    "PostHog query does not return schema columns "
+                            + missingColumns
+                            + ". Alias the selected HogQL columns to match the 
SeaTunnel schema");
+        }
+    }
+
+    private static HttpConnectorException requestFailed(String message) {
+        return new 
HttpConnectorException(HttpConnectorErrorCode.REQUEST_FAILED, message);
+    }
+}
diff --git 
a/seatunnel-connectors-v2/connector-http/connector-http-posthog/src/main/java/org/apache/seatunnel/connectors/seatunnel/posthog/source/config/PostHogSourceOptions.java
 
b/seatunnel-connectors-v2/connector-http/connector-http-posthog/src/main/java/org/apache/seatunnel/connectors/seatunnel/posthog/source/config/PostHogSourceOptions.java
new file mode 100644
index 0000000000..2b0a86be14
--- /dev/null
+++ 
b/seatunnel-connectors-v2/connector-http/connector-http-posthog/src/main/java/org/apache/seatunnel/connectors/seatunnel/posthog/source/config/PostHogSourceOptions.java
@@ -0,0 +1,50 @@
+/*
+ * 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.posthog.source.config;
+
+import org.apache.seatunnel.api.configuration.Option;
+import org.apache.seatunnel.api.configuration.Options;
+
+public final class PostHogSourceOptions {
+
+    private PostHogSourceOptions() {}
+
+    public static final Option<String> BASE_URL =
+            Options.key("base_url")
+                    .stringType()
+                    .defaultValue("https://us.posthog.com";)
+                    .withDescription("PostHog instance base URL");
+
+    public static final Option<String> PROJECT_ID =
+            Options.key("project_id")
+                    .stringType()
+                    .noDefaultValue()
+                    .withDescription("PostHog project ID");
+
+    public static final Option<String> API_KEY =
+            Options.key("api_key")
+                    .stringType()
+                    .noDefaultValue()
+                    .withDescription("PostHog personal API key with query read 
permission");
+
+    public static final Option<String> QUERY =
+            Options.key("query")
+                    .stringType()
+                    .noDefaultValue()
+                    .withDescription("HogQL query to execute");
+}
diff --git 
a/seatunnel-connectors-v2/connector-http/connector-http-posthog/src/main/java/org/apache/seatunnel/connectors/seatunnel/posthog/source/config/PostHogSourceParameter.java
 
b/seatunnel-connectors-v2/connector-http/connector-http-posthog/src/main/java/org/apache/seatunnel/connectors/seatunnel/posthog/source/config/PostHogSourceParameter.java
new file mode 100644
index 0000000000..393a2f5235
--- /dev/null
+++ 
b/seatunnel-connectors-v2/connector-http/connector-http-posthog/src/main/java/org/apache/seatunnel/connectors/seatunnel/posthog/source/config/PostHogSourceParameter.java
@@ -0,0 +1,110 @@
+/*
+ * 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.posthog.source.config;
+
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.common.utils.JsonUtils;
+import org.apache.seatunnel.connectors.seatunnel.http.config.HttpCommonOptions;
+import org.apache.seatunnel.connectors.seatunnel.http.config.HttpParameter;
+import org.apache.seatunnel.connectors.seatunnel.http.config.HttpRequestMethod;
+import org.apache.seatunnel.connectors.seatunnel.http.config.HttpSourceOptions;
+
+import java.io.UnsupportedEncodingException;
+import java.net.URLEncoder;
+import java.util.Collections;
+import java.util.LinkedHashMap;
+import java.util.Map;
+
+public class PostHogSourceParameter extends HttpParameter {
+
+    private static final String AUTHORIZATION = "Authorization";
+    private static final String ACCEPT = "Accept";
+    private static final String CONTENT_TYPE = "Content-Type";
+    private static final String APPLICATION_JSON = "application/json";
+
+    @Override
+    public void buildWithConfig(ReadonlyConfig pluginConfig) {
+        String projectId =
+                
requireNonBlank(pluginConfig.get(PostHogSourceOptions.PROJECT_ID), 
"project_id");
+        String apiKey = 
requireNonBlank(pluginConfig.get(PostHogSourceOptions.API_KEY), "api_key");
+        String query = 
requireNonBlank(pluginConfig.get(PostHogSourceOptions.QUERY), "query");
+        String baseUrl = 
normalizeBaseUrl(pluginConfig.get(PostHogSourceOptions.BASE_URL));
+
+        setUrl(baseUrl + "/api/projects/" + encodePathSegment(projectId) + 
"/query/");
+        setMethod(HttpRequestMethod.POST);
+        setParams(Collections.emptyMap());
+        setKeepParamsAsForm(false);
+
+        Map<String, String> headers =
+                new LinkedHashMap<>(
+                        pluginConfig
+                                .getOptional(HttpCommonOptions.HEADERS)
+                                .orElse(Collections.emptyMap()));
+        setHeader(headers, AUTHORIZATION, "Bearer " + apiKey);
+        setHeader(headers, ACCEPT, APPLICATION_JSON);
+        setHeader(headers, CONTENT_TYPE, APPLICATION_JSON);
+        setHeaders(headers);
+
+        Map<String, Object> queryRequest = new LinkedHashMap<>();
+        queryRequest.put("kind", "HogQLQuery");
+        queryRequest.put("query", query);
+        Map<String, Object> body = new LinkedHashMap<>();
+        body.put("query", queryRequest);
+        body.put("refresh", "blocking");
+        setBody(JsonUtils.toJsonString(body));
+
+        setRetry(pluginConfig.getOptional(HttpCommonOptions.RETRY).orElse(0));
+        setRetryBackoffMultiplierMillis(
+                
pluginConfig.get(HttpCommonOptions.RETRY_BACKOFF_MULTIPLIER_MS));
+        
setRetryBackoffMaxMillis(pluginConfig.get(HttpCommonOptions.RETRY_BACKOFF_MAX_MS));
+        
setConnectTimeoutMs(pluginConfig.get(HttpSourceOptions.CONNECT_TIMEOUT_MS));
+        
setSocketTimeoutMs(pluginConfig.get(HttpSourceOptions.SOCKET_TIMEOUT_MS));
+    }
+
+    private static String normalizeBaseUrl(String baseUrl) {
+        String normalized = requireNonBlank(baseUrl, "base_url");
+        while (normalized.endsWith("/")) {
+            normalized = normalized.substring(0, normalized.length() - 1);
+        }
+        if (normalized.isEmpty()) {
+            throw new IllegalArgumentException("PostHog option 'base_url' must 
not be blank");
+        }
+        return normalized;
+    }
+
+    private static String encodePathSegment(String value) {
+        try {
+            return URLEncoder.encode(value, "UTF-8").replace("+", "%20");
+        } catch (UnsupportedEncodingException e) {
+            throw new IllegalStateException("UTF-8 encoding is not available", 
e);
+        }
+    }
+
+    private static void setHeader(Map<String, String> headers, String name, 
String value) {
+        headers.keySet().removeIf(headerName -> 
headerName.equalsIgnoreCase(name));
+        headers.put(name, value);
+    }
+
+    private static String requireNonBlank(String value, String optionName) {
+        if (value == null || value.trim().isEmpty()) {
+            throw new IllegalArgumentException(
+                    "PostHog option '" + optionName + "' must not be blank");
+        }
+        return value.trim();
+    }
+}
diff --git 
a/seatunnel-connectors-v2/connector-http/connector-http-posthog/src/test/java/org/apache/seatunnel/connectors/seatunnel/posthog/source/PostHogSourceFactoryTest.java
 
b/seatunnel-connectors-v2/connector-http/connector-http-posthog/src/test/java/org/apache/seatunnel/connectors/seatunnel/posthog/source/PostHogSourceFactoryTest.java
new file mode 100644
index 0000000000..eea72d7636
--- /dev/null
+++ 
b/seatunnel-connectors-v2/connector-http/connector-http-posthog/src/test/java/org/apache/seatunnel/connectors/seatunnel/posthog/source/PostHogSourceFactoryTest.java
@@ -0,0 +1,33 @@
+/*
+ * 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.posthog.source;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+public class PostHogSourceFactoryTest {
+
+    @Test
+    public void testFactoryMetadata() {
+        PostHogSourceFactory factory = new PostHogSourceFactory();
+
+        Assertions.assertEquals("PostHog", factory.factoryIdentifier());
+        Assertions.assertEquals(PostHogSource.class, factory.getSourceClass());
+        Assertions.assertNotNull(factory.optionRule());
+    }
+}
diff --git 
a/seatunnel-connectors-v2/connector-http/connector-http-posthog/src/test/java/org/apache/seatunnel/connectors/seatunnel/posthog/source/PostHogSourceReaderTest.java
 
b/seatunnel-connectors-v2/connector-http/connector-http-posthog/src/test/java/org/apache/seatunnel/connectors/seatunnel/posthog/source/PostHogSourceReaderTest.java
new file mode 100644
index 0000000000..ef10f7c9c0
--- /dev/null
+++ 
b/seatunnel-connectors-v2/connector-http/connector-http-posthog/src/test/java/org/apache/seatunnel/connectors/seatunnel/posthog/source/PostHogSourceReaderTest.java
@@ -0,0 +1,294 @@
+/*
+ * 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.posthog.source;
+
+import org.apache.seatunnel.shade.com.fasterxml.jackson.databind.JsonNode;
+
+import org.apache.seatunnel.api.serialization.DeserializationSchema;
+import org.apache.seatunnel.api.source.Collector;
+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.common.utils.JsonUtils;
+import 
org.apache.seatunnel.connectors.seatunnel.common.source.SingleSplitReaderContext;
+import 
org.apache.seatunnel.connectors.seatunnel.http.client.HttpClientProvider;
+import org.apache.seatunnel.connectors.seatunnel.http.client.HttpResponse;
+import org.apache.seatunnel.connectors.seatunnel.http.config.HttpParameter;
+import 
org.apache.seatunnel.connectors.seatunnel.http.exception.HttpConnectorException;
+import 
org.apache.seatunnel.connectors.seatunnel.http.source.DeserializationCollector;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import java.io.IOException;
+import java.nio.charset.StandardCharsets;
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.List;
+import java.util.Map;
+
+public class PostHogSourceReaderTest {
+
+    private PostHogSourceReader reader;
+    private List<SeaTunnelRow> rows;
+    private Collector<SeaTunnelRow> output;
+
+    @BeforeEach
+    public void setUp() {
+        SeaTunnelRowType rowType =
+                new SeaTunnelRowType(
+                        new String[] {"event", "distinct_id"},
+                        new SeaTunnelDataType[] {BasicType.STRING_TYPE, 
BasicType.STRING_TYPE});
+        DeserializationCollector collector =
+                new DeserializationCollector(new 
JsonTextDeserializationSchema());
+        reader = new PostHogSourceReader(new HttpParameter(), null, collector, 
rowType);
+        rows = new ArrayList<>();
+        output = listCollector(rows);
+    }
+
+    @Test
+    public void testMapColumnOrientedResponseToJsonRows() throws Exception {
+        reader.collectResponse(
+                "{\"columns\":[\"distinct_id\",\"event\",\"extra\"],"
+                        + 
"\"results\":[[\"user-1\",\"signup\",{\"plan\":\"pro\"}]],"
+                        + 
"\"query_status\":{\"complete\":true,\"error\":false}}",
+                output);
+
+        Assertions.assertEquals(1, rows.size());
+        JsonNode row = JsonUtils.stringToJsonNode((String) 
rows.get(0).getField(0));
+        Assertions.assertEquals("signup", row.path("event").asText());
+        Assertions.assertEquals("user-1", row.path("distinct_id").asText());
+        Assertions.assertEquals("pro", row.at("/extra/plan").asText());
+    }
+
+    @Test
+    public void testAcceptEmptyResults() throws Exception {
+        
reader.collectResponse("{\"columns\":[\"event\",\"distinct_id\"],\"results\":[]}",
 output);
+
+        Assertions.assertTrue(rows.isEmpty());
+    }
+
+    @Test
+    public void testSendNestedQueryBodyWithoutFlattening() throws Exception {
+        String url = "http://localhost/api/projects/123/query/";;
+        String body =
+                "{\"query\":{\"kind\":\"HogQLQuery\",\"query\":\"SELECT event 
FROM events\"},"
+                        + "\"refresh\":\"blocking\"}";
+        Map<String, String> headers = 
Collections.singletonMap("Authorization", "Bearer phx_test");
+        HttpParameter parameter = new HttpParameter();
+        parameter.setUrl(url);
+        parameter.setHeaders(headers);
+        parameter.setBody(body);
+
+        TestSingleSplitReaderContext context = new 
TestSingleSplitReaderContext();
+        RecordingHttpClientProvider httpClient = new 
RecordingHttpClientProvider();
+        PostHogSourceReader httpReader =
+                new PostHogSourceReader(
+                        parameter,
+                        context,
+                        new DeserializationCollector(new 
JsonTextDeserializationSchema()),
+                        new SeaTunnelRowType(
+                                new String[] {"event", "distinct_id"},
+                                new SeaTunnelDataType[] {
+                                    BasicType.STRING_TYPE, 
BasicType.STRING_TYPE
+                                }));
+        httpReader.setHttpClient(httpClient);
+
+        try {
+            httpReader.internalPollNext(output);
+        } finally {
+            httpReader.close();
+        }
+
+        Assertions.assertEquals(url, httpClient.url);
+        Assertions.assertEquals(headers, httpClient.headers);
+        Assertions.assertEquals(body, httpClient.body);
+        Assertions.assertTrue(context.noMoreElements);
+        Assertions.assertEquals(1, rows.size());
+    }
+
+    @Test
+    public void testRejectMissingSchemaColumn() {
+        HttpConnectorException exception =
+                Assertions.assertThrows(
+                        HttpConnectorException.class,
+                        () ->
+                                reader.collectResponse(
+                                        
"{\"columns\":[\"event\"],\"results\":[[\"signup\"]]}",
+                                        output));
+
+        Assertions.assertTrue(exception.getMessage().contains("distinct_id"));
+        Assertions.assertTrue(exception.getMessage().contains("Alias"));
+    }
+
+    @Test
+    public void testRejectDuplicateColumns() {
+        Assertions.assertThrows(
+                HttpConnectorException.class,
+                () ->
+                        reader.collectResponse(
+                                
"{\"columns\":[\"event\",\"event\"],\"results\":[]}", output));
+    }
+
+    @Test
+    public void testRejectResultWidthMismatch() {
+        Assertions.assertThrows(
+                HttpConnectorException.class,
+                () ->
+                        reader.collectResponse(
+                                "{\"columns\":[\"event\",\"distinct_id\"],"
+                                        + "\"results\":[[\"signup\"]]}",
+                                output));
+    }
+
+    @Test
+    public void testRejectIncompleteBlockingQuery() {
+        Assertions.assertThrows(
+                HttpConnectorException.class,
+                () ->
+                        reader.collectResponse(
+                                
"{\"query_status\":{\"complete\":false,\"error\":false},"
+                                        + "\"columns\":[],\"results\":[]}",
+                                output));
+    }
+
+    @Test
+    public void testRejectQueryError() {
+        HttpConnectorException exception =
+                Assertions.assertThrows(
+                        HttpConnectorException.class,
+                        () ->
+                                reader.collectResponse(
+                                        
"{\"query_status\":{\"complete\":true,\"error\":true,"
+                                                + "\"error_message\":\"invalid 
HogQL\"}}",
+                                        output));
+
+        Assertions.assertTrue(exception.getMessage().contains("invalid 
HogQL"));
+    }
+
+    @Test
+    public void testRejectTopLevelQueryError() {
+        HttpConnectorException exception =
+                Assertions.assertThrows(
+                        HttpConnectorException.class,
+                        () ->
+                                reader.collectResponse(
+                                        "{\"error\":\"query timed 
out\",\"columns\":[],\"results\":[]}",
+                                        output));
+
+        Assertions.assertTrue(exception.getMessage().contains("query timed 
out"));
+    }
+
+    @Test
+    public void testRejectMissingResults() {
+        Assertions.assertThrows(
+                HttpConnectorException.class,
+                () -> 
reader.collectResponse("{\"columns\":[\"event\",\"distinct_id\"]}", output));
+    }
+
+    @Test
+    public void testRejectNonArrayResult() {
+        Assertions.assertThrows(
+                HttpConnectorException.class,
+                () ->
+                        reader.collectResponse(
+                                "{\"columns\":[\"event\",\"distinct_id\"],"
+                                        + 
"\"results\":[{\"event\":\"signup\"}]}",
+                                output));
+    }
+
+    @Test
+    public void testRejectInvalidJson() {
+        Assertions.assertThrows(
+                HttpConnectorException.class, () -> 
reader.collectResponse("not-json", output));
+    }
+
+    private static Collector<SeaTunnelRow> listCollector(List<SeaTunnelRow> 
rows) {
+        return new Collector<SeaTunnelRow>() {
+            @Override
+            public void collect(SeaTunnelRow record) {
+                rows.add(record);
+            }
+
+            @Override
+            public Object getCheckpointLock() {
+                return this;
+            }
+        };
+    }
+
+    private static final class JsonTextDeserializationSchema
+            implements DeserializationSchema<SeaTunnelRow> {
+
+        private final SeaTunnelRowType producedType;
+
+        private JsonTextDeserializationSchema() {
+            this.producedType =
+                    new SeaTunnelRowType(
+                            new String[] {"json"}, new SeaTunnelDataType[] 
{BasicType.STRING_TYPE});
+        }
+
+        @Override
+        public SeaTunnelRow deserialize(byte[] message) throws IOException {
+            return new SeaTunnelRow(new Object[] {new String(message, 
StandardCharsets.UTF_8)});
+        }
+
+        @Override
+        public SeaTunnelDataType<SeaTunnelRow> getProducedType() {
+            return producedType;
+        }
+    }
+
+    private static final class RecordingHttpClientProvider extends 
HttpClientProvider {
+
+        private String url;
+        private Map<String, String> headers;
+        private String body;
+
+        private RecordingHttpClientProvider() {
+            super(new HttpParameter());
+        }
+
+        @Override
+        public HttpResponse doPost(String url, Map<String, String> headers, 
String body) {
+            this.url = url;
+            this.headers = headers;
+            this.body = body;
+            return new HttpResponse(
+                    200,
+                    "{\"columns\":[\"event\",\"distinct_id\"],"
+                            + "\"results\":[[\"signup\",\"user-1\"]]}");
+        }
+    }
+
+    private static final class TestSingleSplitReaderContext extends 
SingleSplitReaderContext {
+
+        private boolean noMoreElements;
+
+        private TestSingleSplitReaderContext() {
+            super(null);
+        }
+
+        @Override
+        public void signalNoMoreElement() {
+            noMoreElements = true;
+        }
+    }
+}
diff --git 
a/seatunnel-connectors-v2/connector-http/connector-http-posthog/src/test/java/org/apache/seatunnel/connectors/seatunnel/posthog/source/config/PostHogSourceParameterTest.java
 
b/seatunnel-connectors-v2/connector-http/connector-http-posthog/src/test/java/org/apache/seatunnel/connectors/seatunnel/posthog/source/config/PostHogSourceParameterTest.java
new file mode 100644
index 0000000000..5687f3984a
--- /dev/null
+++ 
b/seatunnel-connectors-v2/connector-http/connector-http-posthog/src/test/java/org/apache/seatunnel/connectors/seatunnel/posthog/source/config/PostHogSourceParameterTest.java
@@ -0,0 +1,105 @@
+/*
+ * 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.posthog.source.config;
+
+import org.apache.seatunnel.shade.com.fasterxml.jackson.databind.JsonNode;
+
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.common.utils.JsonUtils;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+import java.util.HashMap;
+import java.util.Map;
+
+public class PostHogSourceParameterTest {
+
+    @Test
+    public void testBuildQueryRequest() throws Exception {
+        Map<String, Object> options = new HashMap<>();
+        options.put("base_url", "https://eu.posthog.com/";);
+        options.put("project_id", "project 1/test");
+        options.put("api_key", "phx_test");
+        options.put("query", "SELECT event, distinct_id FROM events LIMIT 10");
+        options.put("retry", 3);
+        options.put("connect_timeout_ms", 1000);
+        options.put("socket_timeout_ms", 2000);
+        Map<String, String> headers = new HashMap<>();
+        headers.put("X-Test", "value");
+        headers.put("authorization", "ignored");
+        headers.put("content-type", "text/plain");
+        options.put("headers", headers);
+
+        PostHogSourceParameter parameter = new PostHogSourceParameter();
+        parameter.buildWithConfig(ReadonlyConfig.fromMap(options));
+
+        Assertions.assertEquals(
+                
"https://eu.posthog.com/api/projects/project%201%2Ftest/query/";,
+                parameter.getUrl());
+        Assertions.assertEquals("post", parameter.getMethod().getMethod());
+        Assertions.assertEquals("Bearer phx_test", 
parameter.getHeaders().get("Authorization"));
+        Assertions.assertEquals("value", parameter.getHeaders().get("X-Test"));
+        Assertions.assertEquals("application/json", 
parameter.getHeaders().get("Accept"));
+        Assertions.assertEquals("application/json", 
parameter.getHeaders().get("Content-Type"));
+        
Assertions.assertFalse(parameter.getHeaders().containsKey("authorization"));
+        
Assertions.assertFalse(parameter.getHeaders().containsKey("content-type"));
+        Assertions.assertEquals(3, parameter.getRetry());
+        Assertions.assertEquals(1000, parameter.getConnectTimeoutMs());
+        Assertions.assertEquals(2000, parameter.getSocketTimeoutMs());
+
+        JsonNode body = JsonUtils.stringToJsonNode(parameter.getBody());
+        Assertions.assertEquals("HogQLQuery", body.at("/query/kind").asText());
+        Assertions.assertEquals(
+                "SELECT event, distinct_id FROM events LIMIT 10", 
body.at("/query/query").asText());
+        Assertions.assertEquals("blocking", body.path("refresh").asText());
+    }
+
+    @Test
+    public void testRejectBlankQuery() {
+        Map<String, Object> options = new HashMap<>();
+        options.put("project_id", "1");
+        options.put("api_key", "phx_test");
+        options.put("query", "  ");
+
+        PostHogSourceParameter parameter = new PostHogSourceParameter();
+        IllegalArgumentException exception =
+                Assertions.assertThrows(
+                        IllegalArgumentException.class,
+                        () -> 
parameter.buildWithConfig(ReadonlyConfig.fromMap(options)));
+
+        Assertions.assertTrue(exception.getMessage().contains("query"));
+    }
+
+    @Test
+    public void testRejectInvalidBaseUrl() {
+        Map<String, Object> options = new HashMap<>();
+        options.put("base_url", "///");
+        options.put("project_id", "1");
+        options.put("api_key", "phx_test");
+        options.put("query", "SELECT event FROM events");
+
+        PostHogSourceParameter parameter = new PostHogSourceParameter();
+        IllegalArgumentException exception =
+                Assertions.assertThrows(
+                        IllegalArgumentException.class,
+                        () -> 
parameter.buildWithConfig(ReadonlyConfig.fromMap(options)));
+
+        Assertions.assertTrue(exception.getMessage().contains("base_url"));
+    }
+}
diff --git a/seatunnel-connectors-v2/connector-http/pom.xml 
b/seatunnel-connectors-v2/connector-http/pom.xml
index a3bddc1c59..7584cef02b 100644
--- a/seatunnel-connectors-v2/connector-http/pom.xml
+++ b/seatunnel-connectors-v2/connector-http/pom.xml
@@ -43,6 +43,7 @@
         <module>connector-http-notion</module>
         <module>connector-http-persistiq</module>
         <module>connector-http-airtable</module>
+        <module>connector-http-posthog</module>
         <module>connector-http-zendesk</module>
     </modules>
 
diff --git a/seatunnel-dist/pom.xml b/seatunnel-dist/pom.xml
index ce4c04f521..0832fe496b 100644
--- a/seatunnel-dist/pom.xml
+++ b/seatunnel-dist/pom.xml
@@ -614,6 +614,12 @@
                     <version>${project.version}</version>
                     <scope>provided</scope>
                 </dependency>
+                <dependency>
+                    <groupId>org.apache.seatunnel</groupId>
+                    <artifactId>connector-http-posthog</artifactId>
+                    <version>${project.version}</version>
+                    <scope>provided</scope>
+                </dependency>
                 <dependency>
                     <groupId>org.apache.seatunnel</groupId>
                     <artifactId>connector-rabbitmq</artifactId>
diff --git 
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-http-e2e/pom.xml 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-http-e2e/pom.xml
index 67b8c17a1a..acea0393d3 100644
--- a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-http-e2e/pom.xml
+++ b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-http-e2e/pom.xml
@@ -123,6 +123,12 @@
             <version>${project.version}</version>
             <scope>test</scope>
         </dependency>
+        <dependency>
+            <groupId>org.apache.seatunnel</groupId>
+            <artifactId>connector-http-posthog</artifactId>
+            <version>${project.version}</version>
+            <scope>test</scope>
+        </dependency>
         <dependency>
             <groupId>org.mock-server</groupId>
             <artifactId>mockserver-netty-no-dependencies</artifactId>
diff --git 
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-http-e2e/src/test/java/org/apache/seatunnel/e2e/connector/http/HttpIT.java
 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-http-e2e/src/test/java/org/apache/seatunnel/e2e/connector/http/HttpIT.java
index 39d351ca05..917844362f 100644
--- 
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-http-e2e/src/test/java/org/apache/seatunnel/e2e/connector/http/HttpIT.java
+++ 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-http-e2e/src/test/java/org/apache/seatunnel/e2e/connector/http/HttpIT.java
@@ -361,6 +361,10 @@ public class HttpIT extends TestSuiteBase implements 
TestResource {
         // http airtable source
         Container.ExecResult execResult22 = 
container.executeJob("/airtable_json_to_assert.conf");
         Assertions.assertEquals(0, execResult22.getExitCode());
+
+        // http posthog source
+        Container.ExecResult execResult23 = 
container.executeJob("/posthog_json_to_assert.conf");
+        Assertions.assertEquals(0, execResult23.getExitCode());
     }
 
     /**
diff --git 
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-http-e2e/src/test/resources/mockserver-config.json
 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-http-e2e/src/test/resources/mockserver-config.json
index 6d7a40b199..b20b25c093 100644
--- 
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-http-e2e/src/test/resources/mockserver-config.json
+++ 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-http-e2e/src/test/resources/mockserver-config.json
@@ -4959,5 +4959,33 @@
         "base64Bytes": "SGVsbG8gV29ybGQgLSBiaW5hcnkgdGVzdCBjb250ZW50"
       }
     }
+  },
+  {
+    "httpRequest": {
+      "method": "POST",
+      "path": "/api/projects/123/query/",
+      "headers": {
+        "Authorization": ["Bearer phx_test"]
+      },
+      "body": {
+        "type": "JSON",
+        "json": "{\"query\":{\"kind\":\"HogQLQuery\",\"query\":\"SELECT event, 
distinct_id FROM events LIMIT 2\"},\"refresh\":\"blocking\"}",
+        "matchType": "ONLY_MATCHING_FIELDS"
+      }
+    },
+    "httpResponse": {
+      "statusCode": 200,
+      "body": {
+        "columns": ["event", "distinct_id"],
+        "results": [
+          ["signup", "user-1"],
+          ["purchase", "user-2"]
+        ],
+        "query_status": {
+          "complete": true,
+          "error": false
+        }
+      }
+    }
   }
-]
\ No newline at end of file
+]
diff --git 
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-http-e2e/src/test/resources/posthog_json_to_assert.conf
 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-http-e2e/src/test/resources/posthog_json_to_assert.conf
new file mode 100644
index 0000000000..cbda2b3218
--- /dev/null
+++ 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-http-e2e/src/test/resources/posthog_json_to_assert.conf
@@ -0,0 +1,67 @@
+#
+# 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 {
+  PostHog {
+    plugin_output = "posthog_events"
+    base_url = "http://mockserver:1080";
+    project_id = "123"
+    api_key = "phx_test"
+    query = "SELECT event, distinct_id FROM events LIMIT 2"
+    schema = {
+      fields {
+        event = string
+        distinct_id = string
+      }
+    }
+  }
+}
+
+sink {
+  Assert {
+    plugin_input = "posthog_events"
+    rules {
+      row_rules = [
+        {
+          rule_type = MIN_ROW
+          rule_value = 2
+        },
+        {
+          rule_type = MAX_ROW
+          rule_value = 2
+        }
+      ]
+      field_rules = [
+        {
+          field_name = event
+          field_type = string
+          field_value = [{ rule_type = NOT_NULL }]
+        },
+        {
+          field_name = distinct_id
+          field_type = string
+          field_value = [{ rule_type = NOT_NULL }]
+        }
+      ]
+    }
+  }
+}

Reply via email to