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-11028-ce4cde54c8abc7a6a9e3831110dbab029512ef7b in repository https://gitbox.apache.org/repos/asf/seatunnel.git
commit b05f17c3179225ba6095185ad80a52a7aaa483cd Author: Gabriel Baldez <[email protected]> AuthorDate: Sat Sep 5 05:17:25 2026 +0000 [Feature][Connector-V2][Shopify] Add Shopify source connector (#11028) Co-authored-by: Daniel Leens <[email protected]> --- .../connectors/changelog/connector-http-shopify.md | 6 + docs/en/connectors/source/Shopify.md | 135 +++++++++++++++++++++ .../connectors/changelog/connector-http-shopify.md | 6 + docs/zh/connectors/source/Shopify.md | 131 ++++++++++++++++++++ plugin-mapping.properties | 1 + .../{ => connector-http-shopify}/pom.xml | 33 ++--- .../seatunnel/shopify/source/ShopifySource.java | 87 +++++++++++++ .../shopify/source/ShopifySourceFactory.java | 64 ++++++++++ .../source/config/ShopifySourceOptions.java | 38 ++++++ .../source/config/ShopifySourceParameter.java | 43 +++++++ .../seatunnel/shopify/ShopifyFactoryTest.java | 31 +++++ .../shopify/ShopifySourceParameterTest.java | 52 ++++++++ .../seatunnel/shopify/ShopifySourceTest.java | 66 ++++++++++ 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 | 57 +++++++++ .../src/test/resources/shopify_json_to_assert.conf | 97 +++++++++++++++ 19 files changed, 842 insertions(+), 22 deletions(-) diff --git a/docs/en/connectors/changelog/connector-http-shopify.md b/docs/en/connectors/changelog/connector-http-shopify.md new file mode 100644 index 0000000000..1cf41e674b --- /dev/null +++ b/docs/en/connectors/changelog/connector-http-shopify.md @@ -0,0 +1,6 @@ +<details><summary> Change Log </summary> + +| Change | Commit | Version | +| --- | --- | --- | + +</details> diff --git a/docs/en/connectors/source/Shopify.md b/docs/en/connectors/source/Shopify.md new file mode 100644 index 0000000000..744f191007 --- /dev/null +++ b/docs/en/connectors/source/Shopify.md @@ -0,0 +1,135 @@ +import ChangeLog from '../changelog/connector-http-shopify.md'; + +# Shopify + +> Shopify source connector + +## Description + +Used to read data from the [Shopify Admin REST API](https://shopify.dev/docs/api/admin-rest). It authenticates with a Shopify Admin API access token (sent in the `X-Shopify-Access-Token` header) and reads a resource endpoint such as orders, products, or customers into SeaTunnel rows. + +## 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) + +## Options + +| name | type | required | default value | +|-----------------------------|---------|----------|---------------| +| url | String | Yes | - | +| access_token | String | Yes | - | +| method | String | No | get | +| headers | Map | No | - | +| schema | Config | No | - | +| format | String | No | text | +| params | Map | No | - | +| body | String | No | - | +| json_field | Config | No | - | +| content_field | String | No | - | +| poll_interval_millis | int | No | - | +| retry | int | No | - | +| retry_backoff_multiplier_ms | int | No | 100 | +| retry_backoff_max_ms | int | No | 10000 | +| json_filed_missed_return_null | boolean | No | false | +| enable_multi_lines | boolean | No | false | +| common-options | config | No | - | + +`pageing` appears in the option rule but is rejected at startup by this connector — see [Pagination](#pagination). + +### url [String] + +The Shopify Admin API endpoint to read from, for example `https://your-store.myshopify.com/admin/api/2024-01/products.json`. + +### access_token [String] + +The Shopify Admin API access token. It is sent in the `X-Shopify-Access-Token` request header. See the [Shopify authentication docs](https://shopify.dev/docs/api/admin-rest#authentication) for how to obtain one. + +### method [String] + +http request method, only supports GET, POST method. + +### headers [Map] + +Extra HTTP request headers. The connector already sets `X-Shopify-Access-Token` from +`access_token` and `Accept: application/json`, so this is only for anything beyond those — +setting `X-Shopify-Access-Token` here is overwritten by `access_token`. + +### schema [Config] + +The structure of the data, including field names and field types. For more details, please refer to [Schema Feature](../../introduction/concepts/schema-feature.md). + +### format [String] + +the format of upstream data, now only support `json` `text`, default `text`. + +### params [Map] + +http params + +### json_field [Config] + +This parameter helps you configure the schema, so this parameter must be used with schema. It maps JSON paths in the response to schema fields. See the [Http source](./Http.md) connector for details and examples. + +### content_field [String] + +This parameter can extract a sub-section of the JSON response (for example the array under a top-level key such as `products` or `orders`) before mapping to rows. See the [Http source](./Http.md) connector for details and examples. + +### common options + +Source plugin common parameters, please refer to [Source Common Options](../common-options/source-common-options.md) for details. + +## Pagination + +**Not supported.** `pageing` is inherited from the HTTP source option rule, but this connector +does not pass it to the reader, so honouring it would read only the first response while the job +reported success. Setting it therefore fails at startup with `HTTP-03`, rather than silently +returning partial data. + +Wiring the inherited pagination through would not help on its own: the shared implementation +reads the next cursor out of the response *body* with a JsonPath, while the Admin REST API +returns it in the `Link` response header. Supporting it properly means teaching +`connector-http-base` to read a header cursor. + +## Example + +```hocon +source { + Shopify { + url = "https://your-store.myshopify.com/admin/api/2024-01/products.json" + access_token = "${SHOPIFY_ACCESS_TOKEN}" + method = "GET" + format = "json" + content_field = "$.products.*" + schema = { + fields { + id = string + title = string + vendor = string + product_type = string + created_at = string + updated_at = string + } + } + } +} +``` + +`${SHOPIFY_ACCESS_TOKEN}` is a SeaTunnel config variable, not an environment variable — it is +substituted only when the value is supplied on the command line: + +```bash +./bin/seatunnel.sh -c your_app.conf -i SHOPIFY_ACCESS_TOKEN=shpat_xxx +``` + +Without `-i`, the literal text `${SHOPIFY_ACCESS_TOKEN}` is sent as the token and Shopify answers +`401`. See [variable configuration](../../introduction/concepts/config.md). + + +## Changelog + +<ChangeLog /> diff --git a/docs/zh/connectors/changelog/connector-http-shopify.md b/docs/zh/connectors/changelog/connector-http-shopify.md new file mode 100644 index 0000000000..1cf41e674b --- /dev/null +++ b/docs/zh/connectors/changelog/connector-http-shopify.md @@ -0,0 +1,6 @@ +<details><summary> Change Log </summary> + +| Change | Commit | Version | +| --- | --- | --- | + +</details> diff --git a/docs/zh/connectors/source/Shopify.md b/docs/zh/connectors/source/Shopify.md new file mode 100644 index 0000000000..adf9181ec5 --- /dev/null +++ b/docs/zh/connectors/source/Shopify.md @@ -0,0 +1,131 @@ +import ChangeLog from '../changelog/connector-http-shopify.md'; + +# Shopify + +> Shopify 数据源连接器 + +## 描述 + +用于从 [Shopify Admin REST API](https://shopify.dev/docs/api/admin-rest) 读取数据。它使用 Shopify Admin API 的访问令牌进行认证(通过 `X-Shopify-Access-Token` 请求头发送),并将某个资源接口(如 orders、products、customers)读取为 SeaTunnel 的行数据。 + +## 主要特性 + +- [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) + +## 选项 + +| 名称 | 类型 | 是否必填 | 默认值 | +|-----------------------------|---------|----------|-----------| +| url | String | 是 | - | +| access_token | String | 是 | - | +| method | String | 否 | get | +| headers | Map | 否 | - | +| schema | Config | 否 | - | +| format | String | 否 | text | +| params | Map | 否 | - | +| body | String | 否 | - | +| json_field | Config | 否 | - | +| content_field | String | 否 | - | +| poll_interval_millis | int | 否 | - | +| retry | int | 否 | - | +| retry_backoff_multiplier_ms | int | 否 | 100 | +| retry_backoff_max_ms | int | 否 | 10000 | +| json_filed_missed_return_null | boolean | 否 | false | +| enable_multi_lines | boolean | 否 | false | +| common-options | config | 否 | - | + +`pageing` 出现在选项规则中,但本连接器会在启动时拒绝它 —— 见[分页](#分页)。 + +### url [String] + +要读取的 Shopify Admin API 接口地址,例如 `https://your-store.myshopify.com/admin/api/2024-01/products.json`。 + +### access_token [String] + +Shopify Admin API 的访问令牌,通过 `X-Shopify-Access-Token` 请求头发送。获取方式请参考 [Shopify 认证文档](https://shopify.dev/docs/api/admin-rest#authentication)。 + +### method [String] + +http 请求方法,仅支持 GET、POST 方法。 + +### headers [Map] + +额外的 HTTP 请求头。连接器已经会根据 `access_token` 设置 `X-Shopify-Access-Token`,并设置 +`Accept: application/json`,因此这里只需要配置这两者之外的请求头;在这里设置 +`X-Shopify-Access-Token` 会被 `access_token` 覆盖。 + +### schema [Config] + +数据的结构,包括字段名称和字段类型。更多详情请参考 [Schema Feature](../../introduction/concepts/schema-feature.md)。 + +### format [String] + +上游数据的格式,目前仅支持 `json` 和 `text`,默认 `text`。 + +### params [Map] + +http 请求参数。 + +### json_field [Config] + +该参数用于配置 schema,因此必须与 schema 一起使用。它将响应中的 JSON 路径映射到 schema 字段。详情和示例请参考 [Http source](./Http.md) 连接器。 + +### content_field [String] + +该参数可以在映射为行之前,提取 JSON 响应中的某个子部分(例如顶层键 `products` 或 `orders` 下的数组)。详情和示例请参考 [Http source](./Http.md) 连接器。 + +### common options + +数据源插件通用参数,详情请参考 [Source Common Options](../common-options/source-common-options.md)。 + +## 分页 + +**不支持。** `pageing` 继承自 HTTP source 的选项规则,但本连接器不会把它传给 reader;若照单 +接受,作业只会读取第一次响应的数据却仍报告成功。因此配置该选项会在启动时以 `HTTP-03` 失败, +而不是静默返回不完整的数据。 + +仅仅把继承来的分页参数传下去也解决不了问题:共享实现是用 JsonPath 从响应*体*中读取下一个游标, +而 Admin REST API 把游标放在 `Link` 响应头里。要真正支持,需要让 `connector-http-base` +能够从响应头读取游标。 + +## 示例 + +```hocon +source { + Shopify { + url = "https://your-store.myshopify.com/admin/api/2024-01/products.json" + access_token = "${SHOPIFY_ACCESS_TOKEN}" + method = "GET" + format = "json" + content_field = "$.products.*" + schema = { + fields { + id = string + title = string + vendor = string + product_type = string + created_at = string + updated_at = string + } + } + } +} +``` + +`${SHOPIFY_ACCESS_TOKEN}` 是 SeaTunnel 的配置变量,不是环境变量 —— 只有在命令行传入对应的值时才会被替换: + +```bash +./bin/seatunnel.sh -c your_app.conf -i SHOPIFY_ACCESS_TOKEN=shpat_xxx +``` + +若不传 `-i`,字面量 `${SHOPIFY_ACCESS_TOKEN}` 会被当作令牌发送,Shopify 将返回 `401`。参见[变量配置](../../introduction/concepts/config.md)。 + + +## 变更日志 + +<ChangeLog /> diff --git a/plugin-mapping.properties b/plugin-mapping.properties index e04731eb67..a340fd6fb9 100644 --- a/plugin-mapping.properties +++ b/plugin-mapping.properties @@ -105,6 +105,7 @@ seatunnel.sink.Tablestore = connector-tablestore seatunnel.source.Tablestore = connector-tablestore seatunnel.source.Lemlist = connector-http-lemlist seatunnel.source.Klaviyo = connector-http-klaviyo +seatunnel.source.Shopify = connector-http-shopify seatunnel.source.Zendesk = connector-http-zendesk seatunnel.source.Stripe = connector-http-stripe seatunnel.sink.Slack = connector-slack diff --git a/seatunnel-connectors-v2/connector-http/pom.xml b/seatunnel-connectors-v2/connector-http/connector-http-shopify/pom.xml similarity index 55% copy from seatunnel-connectors-v2/connector-http/pom.xml copy to seatunnel-connectors-v2/connector-http/connector-http-shopify/pom.xml index 7061a5d53b..5449b4a8a5 100644 --- a/seatunnel-connectors-v2/connector-http/pom.xml +++ b/seatunnel-connectors-v2/connector-http/connector-http-shopify/pom.xml @@ -22,30 +22,19 @@ <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-posthog</module> - <module>connector-http-zendesk</module> - <module>connector-http-stripe</module> - </modules> + <artifactId>connector-http-shopify</artifactId> + <name>SeaTunnel : Connectors V2 : Http : Shopify</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-shopify/src/main/java/org/apache/seatunnel/connectors/seatunnel/shopify/source/ShopifySource.java b/seatunnel-connectors-v2/connector-http/connector-http-shopify/src/main/java/org/apache/seatunnel/connectors/seatunnel/shopify/source/ShopifySource.java new file mode 100644 index 0000000000..3afbdb8504 --- /dev/null +++ b/seatunnel-connectors-v2/connector-http/connector-http-shopify/src/main/java/org/apache/seatunnel/connectors/seatunnel/shopify/source/ShopifySource.java @@ -0,0 +1,87 @@ +/* + * 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.shopify.source; + +import org.apache.seatunnel.api.configuration.ReadonlyConfig; +import org.apache.seatunnel.api.table.type.SeaTunnelRow; +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.config.HttpSourceOptions; +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.HttpSource; +import org.apache.seatunnel.connectors.seatunnel.http.source.HttpSourceReader; +import org.apache.seatunnel.connectors.seatunnel.shopify.source.config.ShopifySourceParameter; + +import lombok.extern.slf4j.Slf4j; + +@Slf4j +/** + * Reads a Shopify Admin REST API resource — orders, products, customers — into SeaTunnel rows. + * + * <p>A thin wrapper over {@link HttpSource}: the only behaviour it adds is authentication, through + * {@link ShopifySourceParameter}, which puts the configured token in {@code + * X-Shopify-Access-Token}. Everything else — schema, format, retries — comes from the HTTP source + * unchanged. + */ +public class ShopifySource extends HttpSource { + private final ShopifySourceParameter shopifySourceParameter = new ShopifySourceParameter(); + + public ShopifySource(ReadonlyConfig pluginConfig) { + super(pluginConfig); + rejectUnsupportedPagination(pluginConfig); + this.shopifySourceParameter.buildWithConfig(pluginConfig); + } + + /** + * {@code pageing} is inherited from the HTTP source option rule, but this connector never hands + * a {@code PageInfo} to its reader, so accepting the option would read the first response and + * report success — a store with more than one page of results would lose the rest silently. + * Failing at startup keeps that from looking like a healthy job. + * + * <p>The Admin REST API returns its cursor in the {@code Link} response header, while the + * shared reader resolves cursors with a JsonPath against the response body, so supporting this + * properly needs a header-cursor mode in {@code connector-http-base} rather than a change here. + */ + private static void rejectUnsupportedPagination(ReadonlyConfig pluginConfig) { + if (pluginConfig.getOptional(HttpSourceOptions.PAGEING).isPresent()) { + throw new HttpConnectorException( + HttpConnectorErrorCode.CONFIG_VALIDATION_FAILED, + "The Shopify source does not support the 'pageing' option: it would be " + + "accepted and then ignored, so only the first page of results would " + + "be read while the job reported success. Remove 'pageing' from the " + + "source configuration."); + } + } + + @Override + public String getPluginName() { + return "Shopify"; + } + + @Override + public AbstractSingleSplitReader<SeaTunnelRow> createReader( + SingleSplitReaderContext readerContext) throws Exception { + return new HttpSourceReader( + this.shopifySourceParameter, + readerContext, + this.deserializationSchema, + jsonField, + contentField); + } +} diff --git a/seatunnel-connectors-v2/connector-http/connector-http-shopify/src/main/java/org/apache/seatunnel/connectors/seatunnel/shopify/source/ShopifySourceFactory.java b/seatunnel-connectors-v2/connector-http/connector-http-shopify/src/main/java/org/apache/seatunnel/connectors/seatunnel/shopify/source/ShopifySourceFactory.java new file mode 100644 index 0000000000..acc7129216 --- /dev/null +++ b/seatunnel-connectors-v2/connector-http/connector-http-shopify/src/main/java/org/apache/seatunnel/connectors/seatunnel/shopify/source/ShopifySourceFactory.java @@ -0,0 +1,64 @@ +/* + * 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.shopify.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.connector.TableSource; +import org.apache.seatunnel.api.table.factory.Factory; +import org.apache.seatunnel.api.table.factory.TableSourceFactoryContext; +import org.apache.seatunnel.connectors.seatunnel.http.source.HttpSourceFactory; +import org.apache.seatunnel.connectors.seatunnel.shopify.source.config.ShopifySourceOptions; + +import com.google.auto.service.AutoService; + +import java.io.Serializable; + +@AutoService(Factory.class) +/** + * Registers {@code Shopify} as a source, discovered through {@code @AutoService} and {@code + * plugin-mapping.properties}. The option rule is the HTTP source's plus a required {@code + * access_token}. + */ +public class ShopifySourceFactory extends HttpSourceFactory { + @Override + public String factoryIdentifier() { + return "Shopify"; + } + + @Override + public <T, SplitT extends SourceSplit, StateT extends Serializable> + TableSource<T, SplitT, StateT> createSource(TableSourceFactoryContext context) { + return () -> (SeaTunnelSource<T, SplitT, StateT>) new ShopifySource(context.getOptions()); + } + + /** + * Without this the factory reports {@code HttpSource.class}, inherited from {@link + * HttpSourceFactory}, so anything reading factory metadata sees the wrong source type. + */ + @Override + public Class<? extends SeaTunnelSource> getSourceClass() { + return ShopifySource.class; + } + + @Override + public OptionRule optionRule() { + return getHttpBuilder().required(ShopifySourceOptions.ACCESS_TOKEN).build(); + } +} diff --git a/seatunnel-connectors-v2/connector-http/connector-http-shopify/src/main/java/org/apache/seatunnel/connectors/seatunnel/shopify/source/config/ShopifySourceOptions.java b/seatunnel-connectors-v2/connector-http/connector-http-shopify/src/main/java/org/apache/seatunnel/connectors/seatunnel/shopify/source/config/ShopifySourceOptions.java new file mode 100644 index 0000000000..2b8b4911e2 --- /dev/null +++ b/seatunnel-connectors-v2/connector-http/connector-http-shopify/src/main/java/org/apache/seatunnel/connectors/seatunnel/shopify/source/config/ShopifySourceOptions.java @@ -0,0 +1,38 @@ +/* + * 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.shopify.source.config; + +import org.apache.seatunnel.api.configuration.Option; +import org.apache.seatunnel.api.configuration.Options; +import org.apache.seatunnel.connectors.seatunnel.http.config.HttpCommonOptions; + +/** + * The options this connector adds to the HTTP source's own: the Admin API access token, and the + * header names used to send it. + */ +public class ShopifySourceOptions extends HttpCommonOptions { + public static final String ACCESS_TOKEN_HEADER = "X-Shopify-Access-Token"; + public static final String ACCEPT = "Accept"; + public static final String APPLICATION_JSON = "application/json"; + + public static final Option<String> ACCESS_TOKEN = + Options.key("access_token") + .stringType() + .noDefaultValue() + .withDescription("Shopify Admin API access token"); +} diff --git a/seatunnel-connectors-v2/connector-http/connector-http-shopify/src/main/java/org/apache/seatunnel/connectors/seatunnel/shopify/source/config/ShopifySourceParameter.java b/seatunnel-connectors-v2/connector-http/connector-http-shopify/src/main/java/org/apache/seatunnel/connectors/seatunnel/shopify/source/config/ShopifySourceParameter.java new file mode 100644 index 0000000000..fabe305cd0 --- /dev/null +++ b/seatunnel-connectors-v2/connector-http/connector-http-shopify/src/main/java/org/apache/seatunnel/connectors/seatunnel/shopify/source/config/ShopifySourceParameter.java @@ -0,0 +1,43 @@ +/* + * 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.shopify.source.config; + +import org.apache.seatunnel.api.configuration.ReadonlyConfig; +import org.apache.seatunnel.connectors.seatunnel.http.config.HttpParameter; + +import java.util.HashMap; + +/** + * Adds Shopify's authentication to the base HTTP parameters: the Admin API expects its token in the + * {@code X-Shopify-Access-Token} header rather than in {@code Authorization}. Any headers the user + * configured are preserved. + */ +public class ShopifySourceParameter extends HttpParameter { + + @Override + public void buildWithConfig(ReadonlyConfig pluginConfig) { + super.buildWithConfig(pluginConfig); + // put the Shopify Admin API access token in headers + this.headers = this.getHeaders() == null ? new HashMap<>() : this.getHeaders(); + this.headers.put(ShopifySourceOptions.ACCEPT, ShopifySourceOptions.APPLICATION_JSON); + this.headers.put( + ShopifySourceOptions.ACCESS_TOKEN_HEADER, + pluginConfig.get(ShopifySourceOptions.ACCESS_TOKEN)); + this.setHeaders(this.headers); + } +} diff --git a/seatunnel-connectors-v2/connector-http/connector-http-shopify/src/test/java/org/apache/seatunnel/connectors/seatunnel/shopify/ShopifyFactoryTest.java b/seatunnel-connectors-v2/connector-http/connector-http-shopify/src/test/java/org/apache/seatunnel/connectors/seatunnel/shopify/ShopifyFactoryTest.java new file mode 100644 index 0000000000..2c49fe4941 --- /dev/null +++ b/seatunnel-connectors-v2/connector-http/connector-http-shopify/src/test/java/org/apache/seatunnel/connectors/seatunnel/shopify/ShopifyFactoryTest.java @@ -0,0 +1,31 @@ +/* + * 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.shopify; + +import org.apache.seatunnel.connectors.seatunnel.shopify.source.ShopifySourceFactory; + +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; + +class ShopifyFactoryTest { + + @Test + void optionRule() { + Assertions.assertNotNull((new ShopifySourceFactory()).optionRule()); + } +} diff --git a/seatunnel-connectors-v2/connector-http/connector-http-shopify/src/test/java/org/apache/seatunnel/connectors/seatunnel/shopify/ShopifySourceParameterTest.java b/seatunnel-connectors-v2/connector-http/connector-http-shopify/src/test/java/org/apache/seatunnel/connectors/seatunnel/shopify/ShopifySourceParameterTest.java new file mode 100644 index 0000000000..d65a7b6f5f --- /dev/null +++ b/seatunnel-connectors-v2/connector-http/connector-http-shopify/src/test/java/org/apache/seatunnel/connectors/seatunnel/shopify/ShopifySourceParameterTest.java @@ -0,0 +1,52 @@ +/* + * 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.shopify; + +import org.apache.seatunnel.api.configuration.ReadonlyConfig; +import org.apache.seatunnel.connectors.seatunnel.shopify.source.config.ShopifySourceParameter; + +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; + +import java.util.HashMap; +import java.util.Map; + +class ShopifySourceParameterTest { + + @Test + void buildWithConfigSetsAccessTokenHeaderAndPreservesExistingHeaders() { + Map<String, String> customHeaders = new HashMap<>(); + customHeaders.put("X-Custom-Header", "custom-value"); + + Map<String, Object> options = new HashMap<>(); + options.put("url", "https://your-store.myshopify.com/admin/api/2024-01/products.json"); + options.put("access_token", "shpat_example_token"); + options.put("headers", customHeaders); + + ShopifySourceParameter parameter = new ShopifySourceParameter(); + parameter.buildWithConfig(ReadonlyConfig.fromMap(options)); + + Map<String, String> headers = parameter.getHeaders(); + // the Shopify Admin API access token is sent in the X-Shopify-Access-Token header + Assertions.assertEquals("shpat_example_token", headers.get("X-Shopify-Access-Token")); + // requests ask for a JSON response + Assertions.assertEquals("application/json", headers.get("Accept")); + // existing custom headers are preserved + Assertions.assertEquals("custom-value", headers.get("X-Custom-Header")); + } +} diff --git a/seatunnel-connectors-v2/connector-http/connector-http-shopify/src/test/java/org/apache/seatunnel/connectors/seatunnel/shopify/ShopifySourceTest.java b/seatunnel-connectors-v2/connector-http/connector-http-shopify/src/test/java/org/apache/seatunnel/connectors/seatunnel/shopify/ShopifySourceTest.java new file mode 100644 index 0000000000..215c470853 --- /dev/null +++ b/seatunnel-connectors-v2/connector-http/connector-http-shopify/src/test/java/org/apache/seatunnel/connectors/seatunnel/shopify/ShopifySourceTest.java @@ -0,0 +1,66 @@ +/* + * 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.shopify; + +import org.apache.seatunnel.api.configuration.ReadonlyConfig; +import org.apache.seatunnel.connectors.seatunnel.http.exception.HttpConnectorException; +import org.apache.seatunnel.connectors.seatunnel.shopify.source.ShopifySource; + +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; + +import java.util.HashMap; +import java.util.Map; + +class ShopifySourceTest { + + private static Map<String, Object> baseOptions() { + Map<String, Object> options = new HashMap<>(); + options.put("url", "https://your-store.myshopify.com/admin/api/2024-01/products.json"); + options.put("access_token", "shpat_example_token"); + return options; + } + + @Test + void aConfigWithoutPagingBuilds() { + Assertions.assertDoesNotThrow( + () -> new ShopifySource(ReadonlyConfig.fromMap(baseOptions()))); + } + + @Test + void pagingIsRejectedRatherThanSilentlyIgnored() { + // The reader never receives a PageInfo, so accepting `pageing` would read the first + // response and report success — a store with more than one page would lose the rest + // without anything in the job saying so. Failing here is the point of the check. + Map<String, Object> paging = new HashMap<>(); + paging.put("page_field", "page"); + paging.put("batch_size", "100"); + + Map<String, Object> options = baseOptions(); + options.put("pageing", paging); + + HttpConnectorException thrown = + Assertions.assertThrows( + HttpConnectorException.class, + () -> new ShopifySource(ReadonlyConfig.fromMap(options))); + + Assertions.assertTrue( + thrown.getMessage().contains("pageing"), + "the message should name the option so the fix is obvious: " + thrown.getMessage()); + } +} diff --git a/seatunnel-connectors-v2/connector-http/pom.xml b/seatunnel-connectors-v2/connector-http/pom.xml index 7061a5d53b..d45f5c4294 100644 --- a/seatunnel-connectors-v2/connector-http/pom.xml +++ b/seatunnel-connectors-v2/connector-http/pom.xml @@ -44,6 +44,7 @@ <module>connector-http-persistiq</module> <module>connector-http-airtable</module> <module>connector-http-posthog</module> + <module>connector-http-shopify</module> <module>connector-http-zendesk</module> <module>connector-http-stripe</module> </modules> diff --git a/seatunnel-dist/pom.xml b/seatunnel-dist/pom.xml index 15156d53f8..11a4681692 100644 --- a/seatunnel-dist/pom.xml +++ b/seatunnel-dist/pom.xml @@ -302,6 +302,12 @@ <version>${project.version}</version> <scope>provided</scope> </dependency> + <dependency> + <groupId>org.apache.seatunnel</groupId> + <artifactId>connector-http-shopify</artifactId> + <version>${project.version}</version> + <scope>provided</scope> + </dependency> <dependency> <groupId>org.apache.seatunnel</groupId> <artifactId>connector-http-zendesk</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 c68e9a14e4..4c5f122512 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 @@ -105,6 +105,12 @@ <version>${project.version}</version> <scope>test</scope> </dependency> + <dependency> + <groupId>org.apache.seatunnel</groupId> + <artifactId>connector-http-shopify</artifactId> + <version>${project.version}</version> + <scope>test</scope> + </dependency> <dependency> <groupId>org.apache.seatunnel</groupId> <artifactId>connector-http-github</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 917844362f..5d6f44ee3a 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 @@ -316,6 +316,10 @@ public class HttpIT extends TestSuiteBase implements TestResource { Container.ExecResult execResult11 = container.executeJob("/persistiq_json_to_assert.conf"); Assertions.assertEquals(0, execResult11.getExitCode()); + // http shopify + Container.ExecResult execResult24 = container.executeJob("/shopify_json_to_assert.conf"); + Assertions.assertEquals(0, execResult24.getExitCode()); + // http httpMultiLine Container.ExecResult execResult12 = container.executeJob("/http_multilinejson_to_assert.conf"); 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 b20b25c093..57c91033b4 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 @@ -3983,6 +3983,63 @@ } } }, + { + "httpRequest": { + "method" : "GET", + "path": "/admin/api/2024-01/products.json", + "headers": { + "X-Shopify-Access-Token": [ + "e2e-shopify-access-token" + ] + } + }, + "httpResponse": { + "body": { + "products": [ + { + "id": "7982391918849", + "title": "Ocean Blue Shirt", + "vendor": "Acme", + "product_type": "Shirts", + "created_at": "2024-01-01T10:15:30-05:00", + "updated_at": "2024-01-02T10:15:30-05:00" + }, + { + "id": "7982391951617", + "title": "Classic Varsity Top", + "vendor": "Acme", + "product_type": "Tops", + "created_at": "2024-01-02T10:15:30-05:00", + "updated_at": "2024-01-03T10:15:30-05:00" + }, + { + "id": "7982391984385", + "title": "Mist Grey Hoodie", + "vendor": "Northwind", + "product_type": "Hoodies", + "created_at": "2024-01-03T10:15:30-05:00", + "updated_at": "2024-01-04T10:15:30-05:00" + }, + { + "id": "7982392017153", + "title": "Draft Denim Jacket", + "vendor": "Northwind", + "product_type": "Jackets", + "created_at": "2024-01-04T10:15:30-05:00", + "updated_at": "2024-01-05T10:15:30-05:00" + }, + { + "id": "7982392049921", + "title": "Pink Ruffle Dress", + "vendor": "Contoso", + "product_type": "Dresses", + "created_at": "2024-01-05T10:15:30-05:00", + "updated_at": "2024-01-06T10:15:30-05:00" + } + ] + } + } + }, { "httpRequest": { "method" : "GET", diff --git a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-http-e2e/src/test/resources/shopify_json_to_assert.conf b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-http-e2e/src/test/resources/shopify_json_to_assert.conf new file mode 100644 index 0000000000..06655b1183 --- /dev/null +++ b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-http-e2e/src/test/resources/shopify_json_to_assert.conf @@ -0,0 +1,97 @@ +# +# 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 { + # Mirrors the configuration documented in docs/en/connectors/source/Shopify.md, including + # content_field: the Admin REST API answers with {"products": [...]} rather than a bare + # array, so the JsonPath is what turns the response into rows. + # + # The mock only answers when X-Shopify-Access-Token is present, so this also covers the one + # thing the connector adds over a plain HTTP source — if buildWithConfig ever stops + # injecting the header, the request does not match and this job fails. + Shopify { + plugin_output = "http" + url = "http://mockserver:1080/admin/api/2024-01/products.json" + access_token = "e2e-shopify-access-token" + method = "GET" + format = "json" + content_field = "$.products.*" + schema = { + fields { + id = string + title = string + vendor = string + product_type = string + created_at = string + updated_at = string + } + } + } +} + +sink { + Assert { + plugin_input = "http" + rules { + row_rules = [ + { + rule_type = MAX_ROW + rule_value = 5 + }, + { + rule_type = MIN_ROW + rule_value = 5 + } + ], + + field_rules = [ + { + field_name = id + field_type = string + field_value = [ + { + rule_type = NOT_NULL + } + ] + }, + { + field_name = title + field_type = string + field_value = [ + { + rule_type = NOT_NULL + } + ] + }, + { + field_name = product_type + field_type = string + field_value = [ + { + rule_type = NOT_NULL + } + ] + } + ] + } + } +}
