diff --git a/airbyte-config/init/src/main/resources/config/STANDARD_SOURCE_DEFINITION/9da77001-af33-4bcd-be46-6252bf9342b9.json b/airbyte-config/init/src/main/resources/config/STANDARD_SOURCE_DEFINITION/9da77001-af33-4bcd-be46-6252bf9342b9.json index bfef7957590f..f22cfd36df8c 100644 --- a/airbyte-config/init/src/main/resources/config/STANDARD_SOURCE_DEFINITION/9da77001-af33-4bcd-be46-6252bf9342b9.json +++ b/airbyte-config/init/src/main/resources/config/STANDARD_SOURCE_DEFINITION/9da77001-af33-4bcd-be46-6252bf9342b9.json @@ -2,6 +2,6 @@ "sourceDefinitionId": "9da77001-af33-4bcd-be46-6252bf9342b9", "name": "Shopify", "dockerRepository": "blotout/source-shopify", - "dockerImageTag": "0.1.28", + "dockerImageTag": "1.43", "documentationUrl": "https://docs.airbyte.io/integrations/sources/shopify" } diff --git a/airbyte-config/init/src/main/resources/seed/source_definitions.yaml b/airbyte-config/init/src/main/resources/seed/source_definitions.yaml index d7726f055bd9..151097d62b48 100644 --- a/airbyte-config/init/src/main/resources/seed/source_definitions.yaml +++ b/airbyte-config/init/src/main/resources/seed/source_definitions.yaml @@ -507,7 +507,7 @@ - name: Shopify sourceDefinitionId: 9da77001-af33-4bcd-be46-6252bf9342b9 dockerRepository: blotout/source-shopify - dockerImageTag: 0.1.28 + dockerImageTag: 1.43 documentationUrl: https://docs.airbyte.io/integrations/sources/shopify sourceType: api - name: Short.io diff --git a/airbyte-config/init/src/main/resources/seed/source_specs.yaml b/airbyte-config/init/src/main/resources/seed/source_specs.yaml index 07b2d1e4c6c5..4c864b488025 100644 --- a/airbyte-config/init/src/main/resources/seed/source_specs.yaml +++ b/airbyte-config/init/src/main/resources/seed/source_specs.yaml @@ -5184,7 +5184,7 @@ supportsNormalization: false supportsDBT: false supported_destination_sync_modes: [] -- dockerImage: "blotout/source-shopify:0.1.28" +- dockerImage: "blotout/source-shopify:1.43" spec: documentationUrl: "https://docs.airbyte.io/integrations/sources/shopify" connectionSpecification: diff --git a/airbyte-integrations/connectors/source-shopify/Dockerfile b/airbyte-integrations/connectors/source-shopify/Dockerfile index f7673702449f..51dd5c831f45 100644 --- a/airbyte-integrations/connectors/source-shopify/Dockerfile +++ b/airbyte-integrations/connectors/source-shopify/Dockerfile @@ -28,5 +28,5 @@ COPY source_shopify ./source_shopify ENV AIRBYTE_ENTRYPOINT "python /airbyte/integration_code/main.py" ENTRYPOINT ["python", "/airbyte/integration_code/main.py"] -LABEL io.airbyte.version=0.1.28 +LABEL io.airbyte.version=1.94 LABEL io.airbyte.name=blotout/source-shopify diff --git a/airbyte-integrations/connectors/source-shopify/acceptance-test-config.yml b/airbyte-integrations/connectors/source-shopify/acceptance-test-config.yml index 045f6bc1c10b..89f5eaa7774e 100644 --- a/airbyte-integrations/connectors/source-shopify/acceptance-test-config.yml +++ b/airbyte-integrations/connectors/source-shopify/acceptance-test-config.yml @@ -25,4 +25,4 @@ tests: full_refresh: - config_path: "secrets/config.json" configured_catalog_path: "integration_tests/configured_catalog.json" - timeout_seconds: 1200 + timeout_seconds: 1200 \ No newline at end of file diff --git a/airbyte-integrations/connectors/source-shopify/integration_tests/abnormal_state.json b/airbyte-integrations/connectors/source-shopify/integration_tests/abnormal_state.json index d4bcc46c5271..0c5d901afaa6 100644 --- a/airbyte-integrations/connectors/source-shopify/integration_tests/abnormal_state.json +++ b/airbyte-integrations/connectors/source-shopify/integration_tests/abnormal_state.json @@ -1,56 +1,81 @@ { "customers": { - "updated_at": "2024-07-19T06:41:50-07:00" + "updated_at": "2025-07-08T05:40:38-07:00" }, "orders": { - "updated_at": "2024-07-19T06:52:06-07:00" + "updated_at": "2025-07-08T05:40:38-07:00" }, "draft_orders": { - "updated_at": "2024-07-07T08:18:59-07:00" + "updated_at": "2025-07-08T05:40:38-07:00" }, "products": { - "updated_at": "2024-07-19T06:56:06-07:00" + "updated_at": "2025-07-08T05:40:38-07:00" }, "abandoned_checkouts": { - "updated_at": "2024-07-08T05:41:48-07:00" + "updated_at": "2025-07-08T05:40:38-07:00" }, - "metafields": { - "updated_at": "2024-07-08T03:38:46-07:00" + "product_variants": { + "id": 99999999999999 }, "collects": { - "id": 99923654213791 + "id": 29427031703741 }, - "custom_collections": { - "updated_at": "2024-07-19T07:01:37-07:00" + "collections": { + "updated_at": "2025-07-08T05:40:38-07:00" }, - "orders_refunds": { - "created_at": "2024-07-19T06:41:47-07:00" + "order_refunds": { + "created_at": "2025-03-03T03:47:46-08:00", + "orders": { + "updated_at": "2025-03-03T03:47:46-08:00" + } }, - "orders_risks": { - "id": 9991307599038 + "order_risks": { + "id": 6446736474301, + "orders": { + "updated_at": "2025-02-22T00:37:28-08:00" + } }, "transactions": { - "created_at": "2024-07-19T06:41:46-07:00" + "created_at": "2025-03-03T03:47:45-08:00", + "orders": { + "updated_at": "2025-03-03T03:47:46-08:00" + } }, "pages": { - "updated_at": "2024-07-08T05:24:11-07:00" + "updated_at": "2025-07-08T05:24:10-07:00" + }, + "metafield_pages": { + "updated_at": "2025-07-08T05:40:38-07:00" }, "price_rules": { - "updated_at": "2024-07-08T05:57:05-07:00" + "updated_at": "2025-09-10T06:48:10-07:00" }, "discount_codes": { - "updated_at": "2024-07-08T05:40:38-07:00" + "updated_at": "2025-09-10T06:48:10-07:00", + "price_rules": { + "updated_at": "2025-09-10T06:48:10-07:00" + } }, - "locations": { - "updated_at": "2024-07-08T05:40:38-07:00" + "inventory_items": { + "updated_at": "2025-02-22T00:40:26-08:00", + "products": { + "updated_at": "2025-08-18T02:39:48-07:00" + } }, "inventory_levels": { - "updated_at": "2024-07-08T05:40:38-07:00" + "updated_at": "2025-03-03T03:47:51-08:00", + "locations": {} }, "fulfillment_orders": { - "id": 9991307599038 + "id": 5424260808893, + "orders": { + "updated_at": "2025-03-03T03:47:46-08:00" + } }, "fulfillments": { - "updated_at": "2024-07-08T05:40:38-07:00" + "updated_at": "2025-02-27T23:49:13-08:00", + "orders": { + "updated_at": "2025-03-03T03:47:46-08:00" + } } -} +} \ No newline at end of file diff --git a/airbyte-integrations/connectors/source-shopify/integration_tests/configured_catalog.json b/airbyte-integrations/connectors/source-shopify/integration_tests/configured_catalog.json index 260bbe53a1a6..83b17154fd38 100644 --- a/airbyte-integrations/connectors/source-shopify/integration_tests/configured_catalog.json +++ b/airbyte-integrations/connectors/source-shopify/integration_tests/configured_catalog.json @@ -98,7 +98,7 @@ }, { "stream": { - "name": "orders_refunds", + "name": "order_refunds", "json_schema": {}, "supported_sync_modes": ["incremental", "full_refresh"], "source_defined_cursor": true, @@ -110,7 +110,7 @@ }, { "stream": { - "name": "orders_risks", + "name": "order_risks", "json_schema": {}, "supported_sync_modes": ["incremental", "full_refresh"], "source_defined_cursor": true, @@ -180,6 +180,18 @@ "cursor_field": ["id"], "destination_sync_mode": "overwrite" }, + { + "stream": { + "name": "inventory_items", + "json_schema": {}, + "supported_sync_modes": ["full_refresh", "incremental"], + "source_defined_cursor": true, + "default_cursor_field": ["updated_at"] + }, + "sync_mode": "incremental", + "cursor_field": ["updated_at"], + "destination_sync_mode": "append" + }, { "stream": { "name": "inventory_levels", @@ -215,6 +227,18 @@ "sync_mode": "incremental", "cursor_field": ["updated_at"], "destination_sync_mode": "append" + }, + { + "stream": { + "name": "product_variants", + "json_schema": {}, + "supported_sync_modes": ["incremental", "full_refresh"], + "source_defined_cursor": true, + "default_cursor_field": ["updated_at"] + }, + "sync_mode": "incremental", + "cursor_field": ["updated_at"], + "destination_sync_mode": "append" } ] } diff --git a/airbyte-integrations/connectors/source-shopify/integration_tests/state.json b/airbyte-integrations/connectors/source-shopify/integration_tests/state.json index 03f98c522caa..0bf5f8bc1458 100644 --- a/airbyte-integrations/connectors/source-shopify/integration_tests/state.json +++ b/airbyte-integrations/connectors/source-shopify/integration_tests/state.json @@ -1,56 +1,76 @@ { - "customers": { - "updated_at": "2021-09-19T09:08:24-07:00" - }, - "orders": { - "updated_at": "2021-09-19T09:08:24-07:00" - }, - "draft_orders": { - "updated_at": "2021-07-07T08:18:58-07:00" - }, - "products": { - "updated_at": "2021-09-19T09:10:43-07:00" - }, - "abandoned_checkouts": { - "updated_at": "2021-07-08T05:41:47-07:00" - }, - "metafields": { - "updated_at": "2021-07-08T03:38:45-07:00" - }, - "collects": { - "id": 29427031703741 - }, - "custom_collections": { - "updated_at": "2021-08-18T02:39:34-07:00" - }, - "orders_refunds": { - "created_at": "2021-09-09T02:57:43-07:00" - }, - "orders_risks": { - "id": 6161307599036 - }, - "transactions": { - "created_at": "2021-09-09T02:57:43-07:00" - }, - "pages": { - "updated_at": "2021-07-08T05:24:10-07:00" - }, - "price_rules": { - "updated_at": "2021-09-10T06:48:10-07:00" - }, - "discount_codes": { - "updated_at": "2021-09-10T06:48:10-07:00" - }, - "locations": { - "updated_at": "2021-09-10T06:48:10-07:00" - }, - "inventory_levels": { - "updated_at": "2021-09-10T06:48:10-07:00" - }, - "fulfillment_orders": { - "id": 123 - }, - "fulfillments": { - "updated_at": "2021-09-10T06:48:10-07:00" - } -} + "customers": { + "updated_at": "2022-06-22T03:50:13-07:00" + }, + "orders": { + "updated_at": "2022-10-10T06:21:53-07:00" + }, + "draft_orders": { + "updated_at": "2022-10-08T05:07:29-07:00" + }, + "products": { + "updated_at": "2022-10-10T06:21:56-07:00" + }, + "abandoned_checkouts": {}, + "metafields": { + "updated_at": "2022-05-30T23:42:02-07:00" + }, + "collects": { + "id": 29427031703740 + }, + "custom_collections": { + "updated_at": "2022-10-08T04:44:51-07:00" + }, + "order_refunds": { + "created_at": "2022-10-10T06:21:53-07:00", + "orders": { + "updated_at": "2022-10-10T06:21:53-07:00" + } + }, + "order_risks": { + "id": 6446736474301, + "orders": { + "updated_at": "2022-03-07T02:09:04-08:00" + } + }, + "transactions": { + "created_at": "2022-10-10T06:21:52-07:00", + "orders": { + "updated_at": "2022-10-10T06:21:53-07:00" + } + }, + "pages": { + "updated_at": "2022-10-08T08:07:00-07:00" + }, + "price_rules": { + "updated_at": "2021-09-10T06:48:10-07:00" + }, + "discount_codes": { + "price_rules": { + "updated_at": "2021-09-10T06:48:10-07:00" + }, + "updated_at": "2021-09-10T06:48:10-07:00" + }, + "inventory_items": { + "products": { + "updated_at": "2022-03-17T03:10:35-07:00" + }, + "updated_at": "2022-03-06T14:12:20-08:00" + }, + "inventory_levels": { + "locations": {}, + "updated_at": "2022-10-10T06:21:56-07:00" + }, + "fulfillment_orders": { + "id": 5567724486845, + "orders": { + "updated_at": "2022-10-10T06:05:29-07:00" + } + }, + "fulfillments": { + "updated_at": "2022-06-22T03:50:14-07:00", + "orders": { + "updated_at": "2022-10-10T06:05:29-07:00" + } + } + } \ No newline at end of file diff --git a/airbyte-integrations/connectors/source-shopify/source_shopify/schemas/customers.json b/airbyte-integrations/connectors/source-shopify/source_shopify/schemas/customers.json index 4a6c9a1b1219..7fba45579c51 100644 --- a/airbyte-integrations/connectors/source-shopify/source_shopify/schemas/customers.json +++ b/airbyte-integrations/connectors/source-shopify/source_shopify/schemas/customers.json @@ -1,5 +1,6 @@ { - "type": "object", + "type": ["null", "object"], + "additionalProperties": true, "properties": { "last_order_name": { "type": ["null", "string"] @@ -13,6 +14,9 @@ "multipass_identifier": { "type": ["null", "string"] }, + "shop_url": { + "type": ["null", "string"] + }, "default_address": { "type": ["null", "object"], "properties": { @@ -69,6 +73,21 @@ } } }, + "email_marketing_consent": { + "type": ["null", "object"], + "properties": { + "consent_updated_at": { + "type": ["null", "string"], + "format": "date-time" + }, + "opt_in_level": { + "type": ["null", "string"] + }, + "state": { + "type": ["null", "string"] + } + } + }, "orders_count": { "type": ["null", "integer"] }, @@ -191,6 +210,27 @@ "created_at": { "type": ["null", "string"], "format": "date-time" + }, + "sms_marketing_consent": { + "type": ["null", "object"], + "properties": { + "consent_updated_at": { + "type": ["null", "string"], + "format": "date-time" + }, + "opt_in_level": { + "type": ["null", "string"] + }, + "state": { + "type": ["null", "string"] + } + } + }, + "tax_exemptions": { + "type": ["null", "string"] + }, + "marketing_opt_in_level": { + "type": ["null", "string"] } } -} +} \ No newline at end of file diff --git a/airbyte-integrations/connectors/source-shopify/source_shopify/schemas/fulfillment_orders.json b/airbyte-integrations/connectors/source-shopify/source_shopify/schemas/fulfillment_orders.json index 16987cbedd70..bbb0e34a038e 100644 --- a/airbyte-integrations/connectors/source-shopify/source_shopify/schemas/fulfillment_orders.json +++ b/airbyte-integrations/connectors/source-shopify/source_shopify/schemas/fulfillment_orders.json @@ -177,4 +177,4 @@ } } } -} +} \ No newline at end of file diff --git a/airbyte-integrations/connectors/source-shopify/source_shopify/schemas/inventory_items.json b/airbyte-integrations/connectors/source-shopify/source_shopify/schemas/inventory_items.json new file mode 100644 index 000000000000..31bf780946bb --- /dev/null +++ b/airbyte-integrations/connectors/source-shopify/source_shopify/schemas/inventory_items.json @@ -0,0 +1,50 @@ +{ + "type": "object", + "additionalProperties": true, + "properties": { + "id": { + "type": ["null", "integer"] + }, + "admin_graphql_api_id": { + "type": ["null", "string"] + }, + "cost": { + "type": ["null", "number"] + }, + "country_code_of_origin": { + "type": ["null", "string"] + }, + "country_harmonized_system_codes": { + "type": ["null", "array"], + "items": { + "type": ["null", "string"] + } + }, + "harmonized_system_code": { + "type": ["null", "string"] + }, + "province_code_of_origin": { + "type": ["null", "string"] + }, + "updated_at": { + "type": ["null", "string"], + "format": "date-time" + }, + "created_at": { + "type": ["null", "string"], + "format": "date-time" + }, + "sku": { + "type": ["null", "string"] + }, + "tracked": { + "type": ["null", "boolean"] + }, + "requires_shipping": { + "type": ["null", "boolean"] + }, + "shop_url": { + "type": ["null", "string"] + } + } +} \ No newline at end of file diff --git a/airbyte-integrations/connectors/source-shopify/source_shopify/schemas/metafields.json b/airbyte-integrations/connectors/source-shopify/source_shopify/schemas/metafields.json index 54102d3951a1..b355b0968477 100644 --- a/airbyte-integrations/connectors/source-shopify/source_shopify/schemas/metafields.json +++ b/airbyte-integrations/connectors/source-shopify/source_shopify/schemas/metafields.json @@ -38,4 +38,4 @@ } }, "type": "object" -} +} \ No newline at end of file diff --git a/airbyte-integrations/connectors/source-shopify/source_shopify/schemas/orders_refunds.json b/airbyte-integrations/connectors/source-shopify/source_shopify/schemas/order_refunds.json similarity index 100% rename from airbyte-integrations/connectors/source-shopify/source_shopify/schemas/orders_refunds.json rename to airbyte-integrations/connectors/source-shopify/source_shopify/schemas/order_refunds.json diff --git a/airbyte-integrations/connectors/source-shopify/source_shopify/schemas/orders_risks.json b/airbyte-integrations/connectors/source-shopify/source_shopify/schemas/order_risks.json similarity index 100% rename from airbyte-integrations/connectors/source-shopify/source_shopify/schemas/orders_risks.json rename to airbyte-integrations/connectors/source-shopify/source_shopify/schemas/order_risks.json diff --git a/airbyte-integrations/connectors/source-shopify/source_shopify/schemas/orders.json b/airbyte-integrations/connectors/source-shopify/source_shopify/schemas/orders.json index dd6194344c13..50930afd2ed9 100644 --- a/airbyte-integrations/connectors/source-shopify/source_shopify/schemas/orders.json +++ b/airbyte-integrations/connectors/source-shopify/source_shopify/schemas/orders.json @@ -1,5 +1,6 @@ { "type": "object", + "additionalProperties": true, "properties": { "id": { "type": ["null", "integer"] @@ -201,12 +202,61 @@ "device_id": { "type": ["null", "string"] }, + "discount_applications": { + "type": ["null", "array"], + "items": { + "type": ["null", "object"], + "properties": { + "type": { + "type": ["null", "string"] + }, + "title": { + "type": ["null", "string"] + }, + "description": { + "type": ["null", "string"] + }, + "value": { + "type": ["null", "string"] + }, + "value_type": { + "type": ["null", "string"] + }, + "allocation_method": { + "type": ["null", "string"] + }, + "target_selection": { + "type": ["null", "string"] + }, + "target_type": { + "type": ["null", "string"] + } + } + } + }, "discount_codes": { - "type": ["null", "array"] + "type": ["null", "array"], + "items": { + "type": ["null", "object"], + "properties": { + "code": { + "type": ["null", "string"] + }, + "amount": { + "type": ["null", "string"] + }, + "type": { + "type": ["null", "string"] + } + } + } }, "email": { "type": ["null", "string"] }, + "estimated_taxes": { + "type": ["null", "boolean"] + }, "financial_status": { "type": ["null", "string"] }, @@ -225,6 +275,9 @@ "location_id": { "type": ["null", "integer"] }, + "merchant_of_record_app_id": { + "type": ["null", "string"] + }, "name": { "type": ["null", "string"] }, @@ -263,6 +316,9 @@ "type": ["null", "string"] } }, + "payment_terms": { + "type": ["null", "string"] + }, "phone": { "type": ["null", "string"] }, @@ -287,10 +343,10 @@ "source_name": { "type": ["null", "string"] }, - "event_name": { + "source_url": { "type": ["null", "string"] }, - "source_url": { + "shop_url": { "type": ["null", "string"] }, "subtotal_price": { @@ -663,9 +719,6 @@ "marketing_opt_in_level": { "type": ["null", "string"] }, - "tax_exemptions": { - "type": ["null", "array"] - }, "admin_graphql_api_id": { "type": ["null", "string"] }, @@ -727,8 +780,59 @@ } } }, - "discount_applications": { - "type": ["null", "array"] + "discount_allocations": { + "type": ["null", "array"], + "items": { + "type": ["null", "object"], + "properties": { + "id": { + "type": ["null", "string"] + }, + "amount": { + "type": ["null", "string"] + }, + "description": { + "type": ["null", "string"] + }, + "created_at": { + "type": ["null", "string"], + "format": "date-time" + }, + "discount_application_index": { + "type": ["null", "number"] + }, + "amount_set": { + "type": ["null", "object"], + "properties": { + "shop_money": { + "type": ["null", "object"], + "properties": { + "amount": { + "type": ["null", "string"] + }, + "currency_code": { + "type": ["null", "string"] + } + } + }, + "presentment_money": { + "type": ["null", "object"], + "properties": { + "amount": { + "type": ["null", "string"] + }, + "currency_code": { + "type": ["null", "string"] + } + } + } + } + }, + "application_type": { + "type": ["null", "string"] + } + } + } }, "fulfillments": { "type": ["null", "array"], @@ -773,13 +877,19 @@ "type": ["null", "string"] }, "tracking_numbers": { - "type": ["null", "array"] + "type": ["null", "array"], + "items": { + "type": ["null", "string"] + } }, "tracking_url": { "type": ["null", "string"] }, "tracking_urls": { - "type": ["null", "array"] + "type": ["null", "array"], + "items": { + "type": ["null", "string"] + } }, "updated_at": { "type": ["null", "string"], @@ -912,7 +1022,18 @@ "type": ["null", "integer"] }, "properties": { - "type": ["null", "array"] + "type": ["null", "array"], + "items": { + "type": ["null", "object"], + "properties": { + "name": { + "type": ["null", "string"] + }, + "value": { + "type": ["null", "string"] + } + } + } }, "quantity": { "type": ["null", "integer"] @@ -1109,9 +1230,19 @@ "items": { "type": ["null", "object"], "properties": { + "id": { + "type": ["null", "string"] + }, "amount": { "type": ["null", "string"] }, + "description": { + "type": ["null", "string"] + }, + "created_at": { + "type": ["null", "string"], + "format": "date-time" + }, "discount_application_index": { "type": ["null", "number"] }, @@ -1141,6 +1272,9 @@ } } } + }, + "application_type": { + "type": ["null", "string"] } } } @@ -1278,7 +1412,18 @@ "type": ["null", "integer"] }, "properties": { - "type": ["null", "array"] + "type": ["null", "array"], + "items": { + "type": ["null", "object"], + "properties": { + "name": { + "type": ["null", "string"] + }, + "value": { + "type": ["null", "string"] + } + } + } }, "quantity": { "type": ["null", "integer"] @@ -1475,9 +1620,19 @@ "items": { "type": ["null", "object"], "properties": { + "id": { + "type": ["null", "string"] + }, "amount": { "type": ["null", "string"] }, + "description": { + "type": ["null", "string"] + }, + "created_at": { + "type": ["null", "string"], + "format": "date-time" + }, "discount_application_index": { "type": ["null", "number"] }, @@ -1507,6 +1662,9 @@ } } } + }, + "application_type": { + "type": ["null", "string"] } } } @@ -1535,59 +1693,6 @@ } }, "refunds": { - "type": ["null", "array"] - }, - "shipping_address": { - "type": ["null", "object"], - "properties": { - "first_name": { - "type": ["null", "string"] - }, - "address1": { - "type": ["null", "string"] - }, - "phone": { - "type": ["null", "string"] - }, - "city": { - "type": ["null", "string"] - }, - "zip": { - "type": ["null", "string"] - }, - "province": { - "type": ["null", "string"] - }, - "country": { - "type": ["null", "string"] - }, - "last_name": { - "type": ["null", "string"] - }, - "address2": { - "type": ["null", "string"] - }, - "company": { - "type": ["null", "string"] - }, - "latitude": { - "type": ["null", "number"] - }, - "longitude": { - "type": ["null", "number"] - }, - "name": { - "type": ["null", "string"] - }, - "country_code": { - "type": ["null", "string"] - }, - "province_code": { - "type": ["null", "string"] - } - } - }, - "shipping_lines": { "type": ["null", "array"], "items": { "type": ["null", "object"], @@ -1595,65 +1700,601 @@ "id": { "type": ["null", "integer"] }, - "carrier_identifier": { + "admin_graphql_api_id": { "type": ["null", "string"] }, - "code": { - "type": ["null", "string"] + "created_at": { + "type": ["null", "string"], + "format": "date-time" }, - "delivery_category": { + "note": { "type": ["null", "string"] }, - "discounted_price": { - "type": ["null", "number"] + "order_id": { + "type": ["null", "integer"] }, - "discounted_price_set": { - "type": ["null", "object"], - "properties": { - "shop_money": { - "type": ["null", "object"], - "properties": { - "amount": { - "type": ["null", "number"] - }, - "currency_code": { - "type": ["null", "string"] - } - } - }, - "presentment_money": { - "type": ["null", "object"], - "properties": { - "amount": { - "type": ["null", "number"] - }, - "currency_code": { - "type": ["null", "string"] - } - } - } - } + "processed_at": { + "type": ["null", "string"], + "format": "date-time" }, - "phone": { - "type": ["null", "string"] + "restock": { + "type": ["null", "boolean"] }, - "price": { - "type": ["null", "number"] + "user_id": { + "type": ["null", "integer"] }, - "price_set": { - "type": ["null", "object"], - "properties": { - "shop_money": { - "type": ["null", "object"], - "properties": { - "amount": { - "type": ["null", "number"] - }, - "currency_code": { - "type": ["null", "string"] - } - } - }, + "order_adjustments": { + "type": ["null", "array"], + "items": { + "type": ["null", "object"], + "properties": { + "amount": { + "type": ["null", "string"] + }, + "amount_set": { + "type": ["null", "object"], + "properties": { + "presentment_money": { + "type": ["null", "object"], + "properties": { + "amount": { + "type": ["null", "string"] + }, + "currency_code": { + "type": ["null", "string"] + } + } + }, + "shop_money": { + "type": ["null", "object"], + "properties": { + "amount": { + "type": ["null", "string"] + }, + "currency_code": { + "type": ["null", "string"] + } + } + } + } + }, + "id": { + "type": ["null", "integer"] + }, + "kind": { + "type": ["null", "string"] + }, + "order_id": { + "type": ["null", "integer"] + }, + "reason": { + "type": ["null", "string"] + }, + "refund_id": { + "type": ["null", "integer"] + }, + "tax_amount": { + "type": ["null", "string"] + }, + "tax_amount_set": { + "type": ["null", "object"], + "properties": { + "presentment_money": { + "type": ["null", "object"], + "properties": { + "amount": { + "type": ["null", "string"] + }, + "currency_code": { + "type": ["null", "string"] + } + } + }, + "shop_money": { + "type": ["null", "object"], + "properties": { + "amount": { + "type": ["null", "string"] + }, + "currency_code": { + "type": ["null", "string"] + } + } + } + } + } + } + } + }, + "transactions": { + "type": ["null", "array"], + "items": { + "type": ["null", "object"], + "properties": { + "id": { + "type": ["null", "integer"] + }, + "admin_graphql_api_id": { + "type": ["null", "string"] + }, + "amount": { + "type": ["null", "string"] + }, + "authorization": { + "type": ["null", "string"] + }, + "created_at": { + "type": ["null", "string"] + }, + "currency": { + "type": ["null", "string"] + }, + "device_id": { + "type": ["null", "integer"] + }, + "error_code": { + "type": ["null", "string"] + }, + "gateway": { + "type": ["null", "string"] + }, + "kind": { + "type": ["null", "string"] + }, + "location_id": { + "type": ["null", "integer"] + }, + "message": { + "type": ["null", "string"] + }, + "order_id": { + "type": ["null", "integer"] + }, + "parent_id": { + "type": ["null", "integer"] + }, + "processed_at": { + "type": ["null", "string"] + }, + "receipt": { + "type": ["null", "object"], + "properties": { + "paid_amount": { + "type": ["null", "string"] + } + } + }, + "source_name": { + "type": ["null", "string"] + }, + "status": { + "type": ["null", "string"] + }, + "test": { + "type": ["null", "boolean"] + }, + "user_id": { + "type": ["null", "integer"] + } + } + } + }, + "refund_line_items": { + "type": ["null", "array"], + "items": { + "type": ["null", "object"], + "properties": { + "id": { + "type": ["null", "integer"] + }, + "line_item_id": { + "type": ["null", "integer"] + }, + "location_id": { + "type": ["null", "integer"] + }, + "quantity": { + "type": ["null", "integer"] + }, + "restock_type": { + "type": ["null", "string"] + }, + "subtotal": { + "type": ["null", "number"] + }, + "subtotal_set": { + "type": ["null", "object"], + "properties": { + "shop_money": { + "type": ["null", "object"], + "properties": { + "amount": { + "type": ["null", "string"] + }, + "currency_code": { + "type": ["null", "string"] + } + } + }, + "presentment_money": { + "type": ["null", "object"], + "properties": { + "amount": { + "type": ["null", "string"] + }, + "currency_code": { + "type": ["null", "string"] + } + } + } + } + }, + "total_tax": { + "type": ["null", "number"] + }, + "total_tax_set": { + "type": ["null", "object"], + "properties": { + "shop_money": { + "type": ["null", "object"], + "properties": { + "amount": { + "type": ["null", "string"] + }, + "currency_code": { + "type": ["null", "string"] + } + } + }, + "presentment_money": { + "type": ["null", "object"], + "properties": { + "amount": { + "type": ["null", "string"] + }, + "currency_code": { + "type": ["null", "string"] + } + } + } + } + }, + "line_item": { + "type": ["null", "object"], + "properties": { + "id": { + "type": ["null", "integer"] + }, + "admin_graphql_api_id": { + "type": ["null", "string"] + }, + "fulfillable_quantity": { + "type": ["null", "integer"] + }, + "fulfillment_service": { + "type": ["null", "string"] + }, + "fulfillment_status": { + "type": ["null", "string"] + }, + "gift_card": { + "type": ["null", "boolean"] + }, + "grams": { + "type": ["null", "number"] + }, + "name": { + "type": ["null", "string"] + }, + "price": { + "type": ["null", "string"] + }, + "price_set": { + "type": ["null", "object"], + "properties": { + "shop_money": { + "type": ["null", "object"], + "properties": { + "amount": { + "type": ["null", "string"] + }, + "currency_code": { + "type": ["null", "string"] + } + } + }, + "presentment_money": { + "type": ["null", "object"], + "properties": { + "amount": { + "type": ["null", "string"] + }, + "currency_code": { + "type": ["null", "string"] + } + } + } + } + }, + "product_exists": { + "type": ["null", "boolean"] + }, + "product_id": { + "type": ["null", "integer"] + }, + "properties": { + "type": ["null", "array"], + "items": { + "type": ["null", "string"] + } + }, + "quantity": { + "type": ["null", "integer"] + }, + "requires_shipping": { + "type": ["null", "boolean"] + }, + "sku": { + "type": ["null", "string"] + }, + "taxable": { + "type": ["null", "boolean"] + }, + "title": { + "type": ["null", "string"] + }, + "total_discount": { + "type": ["null", "string"] + }, + "total_discount_set": { + "type": ["null", "object"], + "properties": { + "shop_money": { + "type": ["null", "object"], + "properties": { + "amount": { + "type": ["null", "string"] + }, + "currency_code": { + "type": ["null", "string"] + } + } + }, + "presentment_money": { + "type": ["null", "object"], + "properties": { + "amount": { + "type": ["null", "string"] + }, + "currency_code": { + "type": ["null", "string"] + } + } + } + } + }, + "variant_id": { + "type": ["null", "integer"] + }, + "variant_inventory_management": { + "type": ["null", "string"] + }, + "variant_title": { + "type": ["null", "string"] + }, + "vendor": { + "type": ["null", "string"] + }, + "tax_lines": { + "type": ["null", "array"], + "items": { + "type": ["null", "object"], + "properties": { + "price": { + "type": ["null", "string"] + }, + "price_set": { + "type": ["null", "object"], + "properties": { + "shop_money": { + "type": ["null", "object"], + "properties": { + "amount": { + "type": ["null", "string"] + }, + "currency_code": { + "type": ["null", "string"] + } + } + }, + "presentment_money": { + "type": ["null", "object"], + "properties": { + "amount": { + "type": ["null", "string"] + }, + "currency_code": { + "type": ["null", "string"] + } + } + } + } + }, + "rate": { + "type": ["null", "number"] + }, + "title": { + "type": ["null", "string"] + } + } + } + }, + "discount_allocations": { + "type": ["null", "array"], + "items": { + "type": ["null", "object"], + "properties": { + "amount": { + "type": ["null", "string"] + }, + "amount_set": { + "type": ["null", "object"], + "properties": { + "shop_money": { + "type": ["null", "object"], + "properties": { + "amount": { + "type": ["null", "string"] + }, + "currency_code": { + "type": ["null", "string"] + } + } + }, + "presentment_money": { + "type": ["null", "object"], + "properties": { + "amount": { + "type": ["null", "string"] + }, + "currency_code": { + "type": ["null", "string"] + } + } + } + } + }, + "discount_application_index": { + "type": ["null", "number"] + } + } + } + } + } + } + } + } + } + } + } + }, + "shipping_address": { + "type": ["null", "object"], + "properties": { + "first_name": { + "type": ["null", "string"] + }, + "address1": { + "type": ["null", "string"] + }, + "phone": { + "type": ["null", "string"] + }, + "city": { + "type": ["null", "string"] + }, + "zip": { + "type": ["null", "string"] + }, + "province": { + "type": ["null", "string"] + }, + "country": { + "type": ["null", "string"] + }, + "last_name": { + "type": ["null", "string"] + }, + "address2": { + "type": ["null", "string"] + }, + "company": { + "type": ["null", "string"] + }, + "latitude": { + "type": ["null", "number"] + }, + "longitude": { + "type": ["null", "number"] + }, + "name": { + "type": ["null", "string"] + }, + "country_code": { + "type": ["null", "string"] + }, + "province_code": { + "type": ["null", "string"] + } + } + }, + "shipping_lines": { + "type": ["null", "array"], + "items": { + "type": ["null", "object"], + "properties": { + "id": { + "type": ["null", "integer"] + }, + "carrier_identifier": { + "type": ["null", "string"] + }, + "code": { + "type": ["null", "string"] + }, + "delivery_category": { + "type": ["null", "string"] + }, + "discounted_price": { + "type": ["null", "number"] + }, + "discounted_price_set": { + "type": ["null", "object"], + "properties": { + "shop_money": { + "type": ["null", "object"], + "properties": { + "amount": { + "type": ["null", "number"] + }, + "currency_code": { + "type": ["null", "string"] + } + } + }, + "presentment_money": { + "type": ["null", "object"], + "properties": { + "amount": { + "type": ["null", "number"] + }, + "currency_code": { + "type": ["null", "string"] + } + } + } + } + }, + "phone": { + "type": ["null", "string"] + }, + "price": { + "type": ["null", "number"] + }, + "price_set": { + "type": ["null", "object"], + "properties": { + "shop_money": { + "type": ["null", "object"], + "properties": { + "amount": { + "type": ["null", "number"] + }, + "currency_code": { + "type": ["null", "string"] + } + } + }, "presentment_money": { "type": ["null", "object"], "properties": { @@ -1680,10 +2321,61 @@ "type": ["null", "array"] }, "discount_allocations": { - "type": ["null", "array"] + "type": ["null", "array"], + "items": { + "type": ["null", "object"], + "properties": { + "id": { + "type": ["null", "string"] + }, + "amount": { + "type": ["null", "string"] + }, + "description": { + "type": ["null", "string"] + }, + "created_at": { + "type": ["null", "string"], + "format": "date-time" + }, + "discount_application_index": { + "type": ["null", "number"] + }, + "amount_set": { + "type": ["null", "object"], + "properties": { + "shop_money": { + "type": ["null", "object"], + "properties": { + "amount": { + "type": ["null", "string"] + }, + "currency_code": { + "type": ["null", "string"] + } + } + }, + "presentment_money": { + "type": ["null", "object"], + "properties": { + "amount": { + "type": ["null", "string"] + }, + "currency_code": { + "type": ["null", "string"] + } + } + } + } + }, + "application_type": { + "type": ["null", "string"] + } + } + } } } } } } -} +} \ No newline at end of file diff --git a/airbyte-integrations/connectors/source-shopify/source_shopify/schemas/product_variants.json b/airbyte-integrations/connectors/source-shopify/source_shopify/schemas/product_variants.json new file mode 100644 index 000000000000..eb28efffe09c --- /dev/null +++ b/airbyte-integrations/connectors/source-shopify/source_shopify/schemas/product_variants.json @@ -0,0 +1,111 @@ +{ + "type": ["null", "object"], + "additionalProperties": true, + "properties": { + "id": { + "type": ["null", "integer"] + }, + "product_id": { + "type": ["null", "integer"] + }, + "title": { + "type": ["null", "string"] + }, + "price": { + "type": ["null", "number"] + }, + "sku": { + "type": ["null", "string"] + }, + "position": { + "type": ["null", "integer"] + }, + "inventory_policy": { + "type": ["null", "string"] + }, + "compare_at_price": { + "type": ["null", "string"] + }, + "fulfillment_service": { + "type": ["null", "string"] + }, + "inventory_management": { + "type": ["null", "string"] + }, + "option1": { + "type": ["null", "string"] + }, + "option2": { + "type": ["null", "string"] + }, + "option3": { + "type": ["null", "string"] + }, + "created_at": { + "type": ["null", "string"], + "format": "date-time" + }, + "updated_at": { + "type": ["null", "string"], + "format": "date-time" + }, + "taxable": { + "type": ["null", "boolean"] + }, + "barcode": { + "type": ["null", "string"] + }, + "grams": { + "type": ["null", "integer"] + }, + "image_id": { + "type": ["null", "integer"] + }, + "weight": { + "type": ["null", "number"] + }, + "weight_unit": { + "type": ["null", "string"] + }, + "inventory_item_id": { + "type": ["null", "integer"] + }, + "inventory_quantity": { + "type": ["null", "integer"] + }, + "old_inventory_quantity": { + "type": ["null", "integer"] + }, + "presentment_prices": { + "type": ["null", "array"], + "items": { + "type": ["null", "object"], + "properties": { + "price": { + "type": ["null", "object"], + "properties": { + "amount": { + "type": ["null", "number"] + }, + "currency_code": { + "type": ["null", "string"] + } + } + }, + "compare_at_price": { + "type": ["null", "number"] + } + } + } + }, + "requires_shipping": { + "type": ["null", "boolean"] + }, + "admin_graphql_api_id": { + "type": ["null", "string"] + }, + "shop_url": { + "type": ["null", "string"] + } + } +} \ No newline at end of file diff --git a/airbyte-integrations/connectors/source-shopify/source_shopify/schemas/products.json b/airbyte-integrations/connectors/source-shopify/source_shopify/schemas/products.json index ac2f37095ef3..fdc190131c46 100644 --- a/airbyte-integrations/connectors/source-shopify/source_shopify/schemas/products.json +++ b/airbyte-integrations/connectors/source-shopify/source_shopify/schemas/products.json @@ -1,5 +1,6 @@ { "type": ["object", "null"], + "additionalProperties": true, "properties": { "published_at": { "type": ["null", "string"], @@ -264,6 +265,9 @@ }, "id": { "type": ["null", "integer"] + }, + "shop_url": { + "type": ["null", "string"] } } -} +} \ No newline at end of file diff --git a/airbyte-integrations/connectors/source-shopify/source_shopify/source.py b/airbyte-integrations/connectors/source-shopify/source_shopify/source.py index ba9675449c0d..b2feef0c370e 100644 --- a/airbyte-integrations/connectors/source-shopify/source_shopify/source.py +++ b/airbyte-integrations/connectors/source-shopify/source_shopify/source.py @@ -4,7 +4,7 @@ from abc import ABC, abstractmethod -from typing import Any, Dict, Iterable, List, Mapping, MutableMapping, Optional, Tuple +from typing import Any, Dict, Iterable, List, Mapping, MutableMapping, Optional, Tuple, Union from urllib.parse import parse_qsl, urlparse import requests @@ -15,6 +15,7 @@ from .auth import ShopifyAuthenticator from .transform import DataTypeEnforcer +from .utils import SCOPES_MAPPING from .utils import EagerlyCachedStreamState as stream_state_cache from .utils import ShopifyRateLimiter as limiter from datetime import datetime @@ -24,16 +25,16 @@ LOGGER = logging.getLogger() class ShopifyStream(HttpStream, ABC): - # Latest Stable Release - api_version = "2021-07" + api_version = "2022-10" # Page size limit = 250 # Define primary key as sort key for full_refresh, or very first sync for incremental_refresh primary_key = "id" order_field = "updated_at" filter_field = "updated_at_min" - filter_field_max = "updated_at_max" + + raise_on_http_errors = True def __init__(self, config: Dict): super().__init__(authenticator=config["authenticator"]) @@ -44,6 +45,13 @@ def __init__(self, config: Dict): def url_base(self) -> str: return f"https://{self.config['shop']}.myshopify.com/admin/api/{self.api_version}/" + @property + def default_filter_field_value(self) -> Union[int, str]: + # certain streams are using `since_id` field as `filter_field`, which requires to use `int` type, + # but many other use `str` values for this, we determine what to use based on `filter_field` value + # by default, we use the user defined `Start Date` as initial value, or 0 for `id`-dependent streams. + return 0 if self.filter_field == "since_id" else self.config["start_date"] + @staticmethod def next_page_token(response: requests.Response) -> Optional[Mapping[str, Any]]: next_page = response.links.get("next", None) @@ -58,22 +66,35 @@ def request_params(self, next_page_token: Mapping[str, Any] = None, **kwargs) -> params.update(**next_page_token) else: params["order"] = f"{self.order_field} asc" - params[self.filter_field] = self.config["start_date"] - - if "no_of_days_to_fetch" in self.config.keys() and self.config["no_of_days_to_fetch"] is not None: - end_date = datetime.strptime(params[self.filter_field], "%Y-%m-%d") + timedelta(days=int(self.config["no_of_days_to_fetch"])) - params[self.filter_field_max] = end_date.strftime("%Y-%m-%d") + params[self.filter_field] = self.default_filter_field_value return params @limiter.balance_rate_limit() def parse_response(self, response: requests.Response, **kwargs) -> Iterable[Mapping]: - json_response = response.json() - records = json_response.get(self.data_field, []) if self.data_field is not None else json_response - # transform method was implemented according to issue 4841 - # Shopify API returns price fields as a string and it should be converted to number - # this solution designed to convert string into number, but in future can be modified for general purpose - for record in records: - yield self._transformer.transform(record) + if response.status_code is requests.codes.OK: + json_response = response.json() + records = json_response.get(self.data_field, []) if self.data_field is not None else json_response + # transform method was implemented according to issue 4841 + # Shopify API returns price fields as a string and it should be converted to number + # this solution designed to convert string into number, but in future can be modified for general purpose + if isinstance(records, dict): + # for cases when we have a single record as dict + # add shop_url to the record to make querying easy + records["shop_url"] = self.config["shop"] + yield self._transformer.transform(records) + else: + # for other cases + for record in records: + # add shop_url to the record to make querying easy + record["shop_url"] = self.config["shop"] + yield self._transformer.transform(record) + + def should_retry(self, response: requests.Response) -> bool: + if response.status_code == 404: + self.logger.warn(f"Stream `{self.name}` is not available, skipping...") + setattr(self, "raise_on_http_errors", False) + return False + return super().should_retry(response) @property @abstractmethod @@ -81,7 +102,6 @@ def data_field(self) -> str: """The name of the field in the response which contains the data""" -# Basic incremental stream class IncrementalShopifyStream(ShopifyStream, ABC): # Setting the check point interval to the limit of the records output @@ -92,8 +112,19 @@ def state_checkpoint_interval(self) -> int: # Setting the default cursor field for all streams cursor_field = "updated_at" + @property + def default_state_comparison_value(self) -> Union[int, str]: + # certain streams are using `id` field as `cursor_field`, which requires to use `int` type, + # but many other use `str` values for this, we determine what to use based on `cursor_field` value + return 0 if self.cursor_field == "id" else "" + def get_updated_state(self, current_stream_state: MutableMapping[str, Any], latest_record: Mapping[str, Any]) -> Mapping[str, Any]: - return {self.cursor_field: max(latest_record.get(self.cursor_field, ""), current_stream_state.get(self.cursor_field, ""))} + return { + self.cursor_field: max( + latest_record.get(self.cursor_field, self.default_state_comparison_value), + current_stream_state.get(self.cursor_field, self.default_state_comparison_value), + ) + } @stream_state_cache.cache_stream_state def request_params(self, stream_state: Mapping[str, Any] = None, next_page_token: Mapping[str, Any] = None, **kwargs): @@ -101,40 +132,179 @@ def request_params(self, stream_state: Mapping[str, Any] = None, next_page_token # If there is a next page token then we should only send pagination-related parameters. if not next_page_token: params["order"] = f"{self.order_field} asc" - - start_date = None if stream_state: params[self.filter_field] = stream_state.get(self.cursor_field) - start_date = datetime.strptime(params[self.filter_field], "%Y-%m-%dT%H:%M:%S%z") - else: - start_date = datetime.strptime(params[self.filter_field], "%Y-%m-%d") - - if "no_of_days_to_fetch" in self.config.keys() and self.config["no_of_days_to_fetch"] is not None: - end_date = start_date + timedelta(days=int(self.config["no_of_days_to_fetch"])) - params[self.filter_field_max] = end_date.strftime("%Y-%m-%d") - return params - # Parse the stream_slice with respect to stream_state for Incremental refresh + # Parse the `stream_slice` with respect to `stream_state` for `Incremental refresh` # cases where we slice the stream, the endpoints for those classes don't accept any other filtering, # but they provide us with the updated_at field in most cases, so we used that as incremental filtering during the order slicing. - def filter_records_newer_than_state(self, stream_state: Mapping[str, Any] = None, records_slice: Mapping[str, Any] = None) -> Iterable: + def filter_records_newer_than_state(self, stream_state: Mapping[str, Any] = None, records_slice: Iterable[Mapping] = None) -> Iterable: # Getting records >= state if stream_state: + state_value = stream_state.get(self.cursor_field) for record in records_slice: - if record.get(self.cursor_field) > stream_state.get(self.cursor_field): + if self.cursor_field in record: + record_value = record.get(self.cursor_field, self.default_state_comparison_value) + if record_value: + if record_value >= state_value: + yield record + else: + # old entities could have cursor field in place, but set to null + self.logger.warn( + f"Stream `{self.name}`, Record ID: `{record.get(self.primary_key)}` cursor value is: {record_value}, record is emitted without state comparison" + ) + yield record + else: + # old entities could miss the cursor field + self.logger.warn( + f"Stream `{self.name}`, Record ID: `{record.get(self.primary_key)}` missing cursor field: {self.cursor_field}, record is emitted without state comparison" + ) yield record else: yield from records_slice +class ShopifySubstream(IncrementalShopifyStream): + """ + ShopifySubstream - provides slicing functionality for streams using parts of data from parent stream. + For example: + - `Refunds Orders` is the entity of `Orders`, + - `OrdersRisks` is the entity of `Orders`, + - `DiscountCodes` is the entity of `PriceRules`, etc. + + :: @ parent_stream - defines the parent stream object to read from + :: @ slice_key - defines the name of the property in stream slices dict. + :: @ nested_record - the name of the field inside of parent stream record. Default is `id`. + :: @ nested_record_field_name - the name of the field inside of nested_record. + :: @ nested_substream - the name of the nested entity inside of parent stream, helps to reduce the number of + API Calls, if present, see `OrderRefunds` stream for more. + """ + + parent_stream_class: object = None + slice_key: str = None + nested_record: str = "id" + nested_record_field_name: str = None + nested_substream = None + nested_substream_list_field_id = None + + @property + def parent_stream(self) -> object: + """ + Returns the instance of parent stream, if the substream has a `parent_stream_class` dependency. + """ + return self.parent_stream_class(self.config) if self.parent_stream_class else None + + def get_updated_state(self, current_stream_state: MutableMapping[str, Any], latest_record: Mapping[str, Any]) -> Mapping[str, Any]: + """UPDATING THE STATE OBJECT: + Stream: Transactions + Parent Stream: Orders + Returns: + { + {...}, + "transactions": { + "created_at": "2022-03-03T03:47:45-08:00", + "orders": { + "updated_at": "2022-03-03T03:47:46-08:00" + } + }, + {...}, + } + """ + updated_state = super().get_updated_state(current_stream_state, latest_record) + # add parent_stream_state to `updated_state` + updated_state[self.parent_stream.name] = stream_state_cache.cached_state.get(self.parent_stream.name) + return updated_state + + def request_params(self, next_page_token: Mapping[str, Any] = None, **kwargs) -> MutableMapping[str, Any]: + params = {"limit": self.limit} + if next_page_token: + params.update(**next_page_token) + return params + + def stream_slices(self, stream_state: Mapping[str, Any] = None, **kwargs) -> Iterable[Optional[Mapping[str, Any]]]: + """ + Reading the parent stream for slices with structure: + EXAMPLE: for given nested_record as `id` of Orders, + + Outputs: + [ + {slice_key: 123}, + {slice_key: 456}, + {...}, + {slice_key: 999 + ] + """ + sorted_substream_slices = [] + + # reading parent nested stream_state from child stream state + parent_stream_state = stream_state.get(self.parent_stream.name) if stream_state else {} + + # reading the parent stream + for record in self.parent_stream.read_records(stream_state=parent_stream_state, **kwargs): + # updating the `stream_state` with the state of it's parent stream + # to have the child stream sync independently from the parent stream + stream_state_cache.cached_state[self.parent_stream.name] = self.parent_stream.get_updated_state({}, record) + # to limit the number of API Calls and reduce the time of data fetch, + # we can pull the ready data for child_substream, if nested data is present, + # and corresponds to the data of child_substream we need. + if self.nested_substream and self.nested_substream_list_field_id: + if record.get(self.nested_substream): + sorted_substream_slices.extend( + [ + { + self.slice_key: sub_record[self.nested_substream_list_field_id], + self.cursor_field: record[self.nested_substream][0].get( + self.cursor_field, self.default_state_comparison_value + ), + } + for sub_record in record[self.nested_record] + ] + ) + elif self.nested_substream: + if record.get(self.nested_substream): + sorted_substream_slices.append( + { + self.slice_key: record[self.nested_record], + self.cursor_field: record[self.nested_substream][0].get(self.cursor_field, self.default_state_comparison_value), + } + ) + else: + yield {self.slice_key: record[self.nested_record]} + + # output slice from sorted list to avoid filtering older records + if self.nested_substream: + if len(sorted_substream_slices) > 0: + # sort by cursor_field + sorted_substream_slices.sort(key=lambda x: x.get(self.cursor_field)) + for sorted_slice in sorted_substream_slices: + yield {self.slice_key: sorted_slice[self.slice_key]} + + def read_records( + self, + stream_state: Mapping[str, Any] = None, + stream_slice: Optional[Mapping[str, Any]] = None, + **kwargs, + ) -> Iterable[Mapping[str, Any]]: + """Reading child streams records for each `id`""" + + slice_data = stream_slice.get(self.slice_key) + # sometimes the stream_slice.get(self.slice_key) has the list of records, + # to avoid data exposition inside the logs, we should get the data we need correctly out of stream_slice. + if isinstance(slice_data, list) and self.nested_record_field_name is not None and len(slice_data) > 0: + slice_data = slice_data[0].get(self.nested_record_field_name) + + self.logger.info(f"Reading {self.name} for {self.slice_key}: {slice_data}") + records = super().read_records(stream_slice=stream_slice, **kwargs) + yield from self.filter_records_newer_than_state(stream_state=stream_state, records_slice=records) + + class Customers(IncrementalShopifyStream): data_field = "customers" def path(self, **kwargs) -> str: return f"{self.data_field}.json" - class Orders(IncrementalShopifyStream): data_field = "orders" @@ -149,15 +319,6 @@ def request_params( params["status"] = "any" return params - def parse_response(self, response: requests.Response, **kwargs) -> Iterable[Mapping]: - - raw_response = response.json()['orders'] - modified_response = [] - for data in raw_response: - data['event_name'] = 'orders' - modified_response.append(data) - return modified_response - class LineItems(Orders): data_field = "orders" fields = "id,created_at,processed_at,updated_at,source_name,line_items,source_name" @@ -198,54 +359,6 @@ def parse_response(self, response: requests.Response, **kwargs) -> Iterable[Mapp pass return modified_response -class ChildSubstream(IncrementalShopifyStream): - - """ - ChildSubstream - provides slicing functionality for streams using parts of data from parent stream. - For example: - - `Refunds Orders` is the entity of `Orders`, - - `OrdersRisks` is the entity of `Orders`, - - `DiscountCodes` is the entity of `PriceRules`, etc. - - :: @ parent_stream_class - defines the parent stream object to read from - :: @ slice_key - defines the name of the property in stream slices dict. - :: @ record_field_name - the name of the field inside of parent stream record. Default is `id`. - """ - - parent_stream_class: object = None - slice_key: str = None - record_field_name: str = "id" - - def request_params(self, next_page_token: Mapping[str, Any] = None, **kwargs) -> MutableMapping[str, Any]: - params = {"limit": self.limit} - if next_page_token: - params.update(**next_page_token) - return params - - def stream_slices(self, stream_state: Mapping[str, Any] = None, **kwargs) -> Iterable[Optional[Mapping[str, Any]]]: - """ - Reading the parent stream for slices with structure: - EXAMPLE: for given record_field_name as `id` of Orders, - - Output: [ {slice_key: 123}, {slice_key: 456}, ..., {slice_key: 999} ] - """ - parent_stream = self.parent_stream_class(self.config) - parent_stream_state = stream_state_cache.cached_state.get(parent_stream.name) - for record in parent_stream.read_records(stream_state=parent_stream_state, **kwargs): - yield {self.slice_key: record[self.record_field_name]} - - def read_records( - self, - stream_state: Mapping[str, Any] = None, - stream_slice: Optional[Mapping[str, Any]] = None, - **kwargs, - ) -> Iterable[Mapping[str, Any]]: - """ Reading child streams records for each `id` """ - - self.logger.info(f"Reading {self.name} for {self.slice_key}: {stream_slice.get(self.slice_key)}") - records = super().read_records(stream_slice=stream_slice, **kwargs) - yield from self.filter_records_newer_than_state(stream_state=stream_state, records_slice=records) - class DraftOrders(IncrementalShopifyStream): data_field = "draft_orders" @@ -253,14 +366,13 @@ class DraftOrders(IncrementalShopifyStream): def path(self, **kwargs) -> str: return f"{self.data_field}.json" - class Products(IncrementalShopifyStream): + use_cache = True data_field = "products" def path(self, **kwargs) -> str: return f"{self.data_field}.json" - class AbandonedCheckouts(IncrementalShopifyStream): data_field = "checkouts" @@ -292,7 +404,6 @@ def path(self, **kwargs) -> str: class Collects(IncrementalShopifyStream): - """ Collects stream does not support Incremental Refresh based on datetime fields, only `since_id` is supported: https://shopify.dev/docs/admin-api/rest/reference/products/collect @@ -300,7 +411,6 @@ class Collects(IncrementalShopifyStream): The Collect stream is the link between Products and Collections, if the Collection is created for Products, the `collect` record is created, it's reasonable to Full Refresh all collects. As for Incremental refresh - we would use the since_id specificaly for this stream. - """ data_field = "collects" @@ -311,37 +421,23 @@ class Collects(IncrementalShopifyStream): def path(self, **kwargs) -> str: return f"{self.data_field}.json" - def get_updated_state(self, current_stream_state: MutableMapping[str, Any], latest_record: Mapping[str, Any]) -> Mapping[str, Any]: - return {self.cursor_field: max(latest_record.get(self.cursor_field, 0), current_stream_state.get(self.cursor_field, 0))} - - def request_params( - self, stream_state: Mapping[str, Any] = None, next_page_token: Mapping[str, Any] = None, **kwargs - ) -> MutableMapping[str, Any]: - params = super().request_params(stream_state=stream_state, next_page_token=next_page_token, **kwargs) - # If there is a next page token then we should only send pagination-related parameters. - if not next_page_token and not stream_state: - params[self.filter_field] = 0 - return params - - -class OrdersRefunds(ChildSubstream): +class OrderRefunds(ShopifySubstream): parent_stream_class: object = Orders slice_key = "order_id" - data_field = "refunds" cursor_field = "created_at" + # we pull out the records that we already know has the refunds data from Orders object + nested_substream = "refunds" def path(self, stream_slice: Mapping[str, Any] = None, **kwargs) -> str: order_id = stream_slice["order_id"] return f"orders/{order_id}/{self.data_field}.json" -class OrdersRisks(ChildSubstream): - +class OrderRisks(ShopifySubstream): parent_stream_class: object = Orders slice_key = "order_id" - data_field = "risks" cursor_field = "id" @@ -349,15 +445,10 @@ def path(self, stream_slice: Mapping[str, Any] = None, **kwargs) -> str: order_id = stream_slice["order_id"] return f"orders/{order_id}/{self.data_field}.json" - def get_updated_state(self, current_stream_state: MutableMapping[str, Any], latest_record: Mapping[str, Any]) -> Mapping[str, Any]: - return {self.cursor_field: max(latest_record.get(self.cursor_field, 0), current_stream_state.get(self.cursor_field, 0))} - - -class Transactions(ChildSubstream): +class Transactions(ShopifySubstream): parent_stream_class: object = Orders slice_key = "order_id" - data_field = "transactions" cursor_field = "created_at" @@ -379,12 +470,22 @@ class PriceRules(IncrementalShopifyStream): def path(self, **kwargs) -> str: return f"{self.data_field}.json" +class ProductVariants(ShopifySubstream): + parent_stream_class: object = Products + cursor_field = "id" + slice_key = "product_id" + data_field = "variants" + nested_substream = "variants" + filter_field = "since_id" + + def path(self, stream_slice: Mapping[str, Any] = None, **kwargs) -> str: + product_id = stream_slice[self.slice_key] + return f"products/{product_id}/{self.data_field}.json" -class DiscountCodes(ChildSubstream): +class DiscountCodes(ShopifySubstream): parent_stream_class: object = PriceRules slice_key = "price_rule_id" - data_field = "discount_codes" def path(self, stream_slice: Mapping[str, Any] = None, **kwargs) -> str: @@ -393,7 +494,6 @@ def path(self, stream_slice: Mapping[str, Any] = None, **kwargs) -> str: class Locations(ShopifyStream): - """ The location API does not support any form of filtering. https://shopify.dev/api/admin-rest/2021-07/resources/location @@ -407,11 +507,9 @@ def path(self, **kwargs): return f"{self.data_field}.json" -class InventoryLevels(ChildSubstream): +class InventoryLevels(ShopifySubstream): parent_stream_class: object = Locations slice_key = "location_id" - cursor_field = "updated_at" - data_field = "inventory_levels" def path(self, stream_slice: Mapping[str, Any] = None, **kwargs) -> str: @@ -426,34 +524,35 @@ def generate_key(record): return record # associate the surrogate key - yield from map( - generate_key, - records_stream, - ) + yield from map(generate_key, records_stream) + +class InventoryItems(ShopifySubstream): + parent_stream_class: object = Products + slice_key = "id" + nested_record = "variants" + nested_record_field_name = "inventory_item_id" + data_field = "inventory_items" + + def path(self, stream_slice: Mapping[str, Any] = None, **kwargs) -> str: + ids = ",".join(str(x[self.nested_record_field_name]) for x in stream_slice[self.slice_key]) + return f"inventory_items.json?ids={ids}" -class FulfillmentOrders(ChildSubstream): +class FulfillmentOrders(ShopifySubstream): parent_stream_class: object = Orders slice_key = "order_id" - data_field = "fulfillment_orders" - cursor_field = "id" def path(self, stream_slice: Mapping[str, Any] = None, **kwargs) -> str: order_id = stream_slice[self.slice_key] return f"orders/{order_id}/{self.data_field}.json" - def get_updated_state(self, current_stream_state: MutableMapping[str, Any], latest_record: Mapping[str, Any]) -> Mapping[str, Any]: - return {self.cursor_field: max(latest_record.get(self.cursor_field, 0), current_stream_state.get(self.cursor_field, 0))} - - -class Fulfillments(ChildSubstream): +class Fulfillments(ShopifySubstream): parent_stream_class: object = Orders slice_key = "order_id" - data_field = "fulfillments" def path(self, stream_slice: Mapping[str, Any] = None, **kwargs) -> str: @@ -468,7 +567,7 @@ def check_connection(self, logger: AirbyteLogger, config: Mapping[str, Any]) -> Testing connection availability for the connector. """ auth = ShopifyAuthenticator(config).get_auth_header() - api_version = "2021-07" # Latest Stable Release + api_version = "2021-10" # Latest Stable Release url = f"https://{config['shop']}.myshopify.com/admin/api/{api_version}/shop.json" try: @@ -494,17 +593,19 @@ def streams(self, config: Mapping[str, Any]) -> List[Stream]: LineItems(config), DraftOrders(config), Products(config), + ProductVariants(config), AbandonedCheckouts(config), Metafields(config), CustomCollections(config), Collects(config), - OrdersRefunds(config), - OrdersRisks(config), + OrderRefunds(config), + OrderRisks(config), Transactions(config), Pages(config), PriceRules(config), DiscountCodes(config), Locations(config), + InventoryItems(config), InventoryLevels(config), FulfillmentOrders(config), Fulfillments(config), diff --git a/airbyte-integrations/connectors/source-shopify/source_shopify/utils.py b/airbyte-integrations/connectors/source-shopify/source_shopify/utils.py index d0c5e7d4b540..884911b83c3b 100644 --- a/airbyte-integrations/connectors/source-shopify/source_shopify/utils.py +++ b/airbyte-integrations/connectors/source-shopify/source_shopify/utils.py @@ -1,5 +1,5 @@ # -# Copyright (c) 2021 Airbyte, Inc., all rights reserved. +# Copyright (c) 2022 Airbyte, Inc., all rights reserved. # @@ -9,9 +9,45 @@ import requests +SCOPES_MAPPING = { + "read_customers": ["Customers", "MetafieldCustomers"], + "read_orders": [ + "Orders", + "AbandonedCheckouts", + "TenderTransactions", + "Transactions", + "Fulfillments", + "OrderRefunds", + "OrderRisks", + "MetafieldOrders", + ], + "read_draft_orders": ["DraftOrders", "MetafieldDraftOrders"], + "read_products": [ + "Products", + "MetafieldProducts", + "ProductImages", + "MetafieldProductImages", + "MetafieldProductVariants", + "CustomCollections", + "Collects", + "Collections", + "ProductVariants", + "MetafieldCollections", + "SmartCollections", + "MetafieldSmartCollections", + ], + "read_content": ["Pages", "MetafieldPages"], + "read_price_rules": ["PriceRules"], + "read_discounts": ["DiscountCodes"], + "read_locations": ["Locations", "MetafieldLocations"], + "read_inventory": ["InventoryItems", "InventoryLevels"], + "read_merchant_managed_fulfillment_orders": ["FulfillmentOrders"], + "read_shopify_payments_payouts": ["BalanceTransactions"], + "read_online_store_pages": ["Articles", "MetafieldArticles", "Blogs", "MetafieldBlogs"], +} -class ShopifyRateLimiter: +class ShopifyRateLimiter: """ Define timings for RateLimits. Adjust timings if needed. @@ -92,7 +128,6 @@ def wrapper_balance_rate_limit(*args, **kwargs): class EagerlyCachedStreamState: - """ This is the placeholder for the tmp stream state for each incremental stream, It's empty, once the sync has started and is being updated while sync operation takes place, @@ -127,6 +162,7 @@ def stream_state_to_tmp(*args, state_object: Dict = cached_state, **kwargs) -> D state_object[stream.name] = { stream.cursor_field: min(current_stream_state.get(stream.cursor_field, ""), tmp_stream_state_value) } + return state_object def cache_stream_state(func): @@ -135,4 +171,4 @@ def decorator(*args, **kwargs): EagerlyCachedStreamState.stream_state_to_tmp(*args, **kwargs) return func(*args, **kwargs) - return decorator + return decorator \ No newline at end of file