Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
21 commits
Select commit Hold shift + click to select a range
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 6 additions & 4 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
@@ -1,5 +1,11 @@
# Changelog

## [Version 1.6.1](https://github.com/dataiku/dss-plugin-api-connect/releases/tag/v1.6.1) - Feature and bugfix - 2026-07-27

- Fix templating for multiform body
- Adding a configurable retry for several HTTP errors
- Fix issue with empty cell(s) in bigint columns of the recipe's input dataset

## Version 1.6.0 - Enhancement release - 2026-08-20

- Added supported Python versions: 3.12, 3.13, 3.14
Expand All @@ -10,10 +16,6 @@

- Added Cobuild support to the API Connect recipe

## [Version 1.4.3](https://github.com/dataiku/dss-plugin-api-connect/releases/tag/v1.4.3) - Bugfix - 2026-07-27

- Fix templating for multiform body

## [Version 1.4.2](https://github.com/dataiku/dss-plugin-api-connect/releases/tag/v1.4.2) - Bugfix - 2026-07-22

- Remove user's 'Other credentials' from recipe's logs
Expand Down
71 changes: 71 additions & 0 deletions custom-recipes/api-connect/recipe.json
Original file line number Diff line number Diff line change
Expand Up @@ -411,6 +411,77 @@
"description": "-1 for no limit",
"type": "INT",
"defaultValue": -1
},
{
"name": "http_errors_retry_strategy",
"label": "Retry on error logic",
"description": "",
"type": "SELECT",
"defaultValue": null,
"selectChoices":[
{"value": null, "label": "No retry"},
{"value": "linear", "label": "Linear backoff"},
{"value": "exponential", "label": "Exponential backoff"}
]
},
{
"name": "http_errors_to_retry",
"label": "Errors to retry",
"description": "Click to select errors that can trigger a retry",
"type": "MULTISELECT",
"defaultValue": null,
"selectChoices":[
{"value": "408", "label": "408 Request Timeout"},
{"value": "429", "label": "429 Too many requests"},
{"value": "502", "label": "502 Bad Gateway"},
{"value": "503", "label": "503 Service Unavailable"},
{"value": "504", "label": "504 Gateway Time out"}
],
"visibilityCondition": "(['exponential', 'linear'].includes(model.http_errors_retry_strategy))"
},
{
"name": "http_errors_retry_scope",
"label": "Retry scope",
"description": "Apply the retry budget to the entire input dataset or independently to each input row",
"type": "SELECT",
"defaultValue": "dataset",
"selectChoices":[
{"value": "dataset", "label": "Per dataset"},
{"value": "row", "label": "Per row"}
],
"visibilityCondition": "(['exponential', 'linear'].includes(model.http_errors_retry_strategy))"
},
{
"name": "http_errors_initial_delay",
"label": "Initial delay",
"description": "in seconds",
"type": "INT",
"defaultValue": 1,
"visibilityCondition": "(['exponential'].includes(model.http_errors_retry_strategy))"
},
{
"name": "http_errors_maximum_delay",
"label": "Maximum delay",
"description": "in seconds",
"type": "INT",
"defaultValue": 120,
"visibilityCondition": "(['exponential'].includes(model.http_errors_retry_strategy))"
},
{
"name": "http_errors_delay",
"label": "Delay",
"description": "in seconds",
"type": "INT",
"defaultValue": 1,
"visibilityCondition": "(['linear'].includes(model.http_errors_retry_strategy))"
},
{
"name": "http_errors_maximum_retries",
"label": "Maximum number of retries",
"description": "Number of times to retry a request after an error",
"type": "INT",
"defaultValue": 5,
"visibilityCondition": "(['linear'].includes(model.http_errors_retry_strategy))"
}
],
"resourceKeys": []
Expand Down
23 changes: 17 additions & 6 deletions custom-recipes/api-connect/recipe.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,10 +2,11 @@
import dataiku
from dataiku.customrecipe import get_input_names_for_role, get_recipe_config, get_output_names_for_role
import pandas as pd
from safe_logger import SafeLogger
from dku_utils import get_dku_key_values, get_endpoint_parameters, get_secure_credentials, get_user_secrets
from rest_api_recipe_session import RestApiRecipeSession
from dku_constants import DKUConstants
from api_connect_safe_logger import SafeLogger
from api_connect_dku_utils import get_dku_key_values, get_endpoint_parameters, get_secure_credentials, get_user_secrets, get_retry_handler_parameters_from_config
from api_connect_rest_api_recipe_session import RestApiRecipeSession
from api_connect_dku_constants import DKUConstants
from api_connect_retry_handler import RetryHandler


logger = SafeLogger("api-connect plugin", forbidden_keys=DKUConstants.FORBIDDEN_KEYS)
Expand Down Expand Up @@ -49,10 +50,18 @@ def get_partitioning_keys(id_list, dku_flow_variables):
custom_key_values.update(user_secrets)
display_metadata = config.get("display_metadata", False)
maximum_number_rows = config.get("maximum_number_rows", -1)
retry_scope = config.get("http_errors_retry_scope", "dataset")
input_parameters_dataset = dataiku.Dataset(input_A_names[0])
partitioning_keys = get_partitioning_keys(input_parameters_dataset, dku_flow_variables)
custom_key_values.update(partitioning_keys)
input_parameters_dataframe = input_parameters_dataset.get_dataframe(infer_with_pandas=False)
input_parameters_dataframe = input_parameters_dataset.get_dataframe(infer_with_pandas=False, use_nullable_integers=True)
backoff_type, initial_delay, maximum_number_of_retries, maximum_duration_of_retry, status_codes_to_retry = get_retry_handler_parameters_from_config(config)
retry_handler = None
if backoff_type:
retry_handler = RetryHandler(
backoff_type=backoff_type, initial_delay=initial_delay, maximum_number_of_retries=maximum_number_of_retries,
maximum_duration_of_retry=maximum_duration_of_retry, status_codes_to_retry=status_codes_to_retry
)

recipe_session = RestApiRecipeSession(
custom_key_values,
Expand All @@ -64,7 +73,9 @@ def get_partitioning_keys(id_list, dku_flow_variables):
parameter_renamings,
display_metadata,
maximum_number_rows=maximum_number_rows,
behaviour_when_error=behaviour_when_error
behaviour_when_error=behaviour_when_error,
retry_handler=retry_handler,
retry_scope=retry_scope
)
results = recipe_session.process_dataframe(input_parameters_dataframe, is_raw_output)

Expand Down
59 changes: 59 additions & 0 deletions python-connectors/api-connect_dataset/connector.json
Original file line number Diff line number Diff line change
Expand Up @@ -356,6 +356,65 @@
"description": "-1 for no limit",
"type": "INT",
"defaultValue": -1
},
{
"name": "http_errors_retry_strategy",
"label": "Retry on error logic",
"description": "",
"type": "SELECT",
"defaultValue": null,
"selectChoices":[
{"value": null, "label": "No retry"},
{"value": "linear", "label": "Linear backoff"},
{"value": "exponential", "label": "Exponential backoff"}
]
},
{
"name": "http_errors_to_retry",
"label": "Errors to retry",
"description": "Click to select errors that can trigger a retry",
"type": "MULTISELECT",
"defaultValue": null,
"selectChoices":[
{"value": "408", "label": "408 Request Timeout"},
{"value": "429", "label": "429 Too many requests"},
{"value": "502", "label": "502 Bad Gateway"},
{"value": "503", "label": "503 Service Unavailable"},
{"value": "504", "label": "504 Gateway Time out"}
],
"visibilityCondition": "(['exponential', 'linear'].includes(model.http_errors_retry_strategy))"
},
{
"name": "http_errors_initial_delay",
"label": "Initial delay",
"description": "in seconds",
"type": "INT",
"defaultValue": 1,
"visibilityCondition": "(['exponential'].includes(model.http_errors_retry_strategy))"
},
{
"name": "http_errors_maximum_delay",
"label": "Maximum delay",
"description": "in seconds",
"type": "INT",
"defaultValue": 120,
"visibilityCondition": "(['exponential'].includes(model.http_errors_retry_strategy))"
},
{
"name": "http_errors_delay",
"label": "Delay",
"description": "in seconds",
"type": "INT",
"defaultValue": 1,
"visibilityCondition": "(['linear'].includes(model.http_errors_retry_strategy))"
},
{
"name": "http_errors_maximum_retries",
"label": "Maximum number of retries",
"description": "Number of times to retry a request after an error",
"type": "INT",
"defaultValue": 5,
"visibilityCondition": "(['linear'].includes(model.http_errors_retry_strategy))"
}
]
}
20 changes: 14 additions & 6 deletions python-connectors/api-connect_dataset/connector.py
Original file line number Diff line number Diff line change
@@ -1,14 +1,15 @@
from dataiku.connector import Connector
from dataikuapi.utils import DataikuException
from safe_logger import SafeLogger
from rest_api_client import RestAPIClient
from dku_utils import (
from api_connect_safe_logger import SafeLogger
from api_connect_rest_api_client import RestAPIClient
from api_connect_dku_utils import (
get_dku_key_values, get_endpoint_parameters,
parse_keys_for_json, get_value_from_path, get_secure_credentials,
decode_csv_data, decode_bytes, get_user_secrets
decode_csv_data, decode_bytes, get_user_secrets, get_retry_handler_parameters_from_config
)
from dku_constants import DKUConstants
from api_connect_dku_constants import DKUConstants
import json
from api_connect_retry_handler import RetryHandler


logger = SafeLogger("api-connect plugin", forbidden_keys=DKUConstants.FORBIDDEN_KEYS)
Expand All @@ -26,7 +27,14 @@ def __init__(self, config, plugin_config):
custom_key_values = get_dku_key_values(config.get("custom_key_values", {}))
user_secrets = get_user_secrets(config)
custom_key_values.update(user_secrets)
self.client = RestAPIClient(credential, secure_credentials, endpoint_parameters, custom_key_values)
backoff_type, initial_delay, maximum_number_of_retries, maximum_duration_of_retry, status_codes_to_retry = get_retry_handler_parameters_from_config(config)
retry_handler = None
if backoff_type:
retry_handler = RetryHandler(
backoff_type=backoff_type, initial_delay=initial_delay, maximum_number_of_retries=maximum_number_of_retries,
maximum_duration_of_retry=maximum_duration_of_retry, status_codes_to_retry=status_codes_to_retry
)
self.client = RestAPIClient(credential, secure_credentials, endpoint_parameters, custom_key_values, retry_handler=retry_handler)
extraction_key = endpoint_parameters.get("extraction_key", None)
self.extraction_key = extraction_key or ''
self.extraction_path = self.extraction_key.split('.')
Expand Down
File renamed without changes.
20 changes: 19 additions & 1 deletion python-lib/dku_utils.py → python-lib/api_connect_dku_utils.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,7 @@
import math
from collections import defaultdict
from jsonpath_ng.ext import parse
from safe_logger import SafeLogger
from api_connect_safe_logger import SafeLogger


logger = SafeLogger("api-connect plugin utils")
Expand Down Expand Up @@ -319,3 +319,21 @@ def join_url(base_url, segment):
segment = segment.lstrip("/")
segments.append(segment)
return "/".join(segments)


def get_retry_handler_parameters_from_config(config):
backoff_type = initial_delay = maximum_number_of_retries = maximum_duration_of_retry = status_codes_to_retry = None
http_errors_retry_strategy = config.get("http_errors_retry_strategy", None)
if not http_errors_retry_strategy:
return backoff_type, initial_delay, maximum_number_of_retries, maximum_duration_of_retry, status_codes_to_retry
if http_errors_retry_strategy in ["linear", "exponential"]:
backoff_type = http_errors_retry_strategy
if backoff_type == "linear":
initial_delay = config.get("http_errors_delay")
maximum_number_of_retries = config.get("http_errors_maximum_retries", None)
if backoff_type == "exponential":
initial_delay = config.get("http_errors_initial_delay")
maximum_duration_of_retry = config.get("http_errors_maximum_delay", None)
if backoff_type:
status_codes_to_retry = config.get("http_errors_to_retry", [])
return backoff_type, initial_delay, maximum_number_of_retries, maximum_duration_of_retry, status_codes_to_retry
File renamed without changes.
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
from safe_logger import SafeLogger
from dku_utils import get_value_from_path, extract_key_using_json_path, join_url
from api_connect_safe_logger import SafeLogger
from api_connect_dku_utils import get_value_from_path, extract_key_using_json_path, join_url


logger = SafeLogger("api-connect plugin Pagination")
Expand Down
File renamed without changes.
Original file line number Diff line number Diff line change
Expand Up @@ -2,12 +2,13 @@
import time
import copy
import tempfile
from pagination import Pagination
from safe_logger import SafeLogger
from loop_detector import LoopDetector
from dku_utils import get_dku_key_values, get_dku_duplicated_key_values, template_dict, format_template, is_reponse_xml, xml_to_json
from dku_constants import DKUConstants
from rest_api_auth import get_auth
from api_connect_pagination import Pagination
from api_connect_safe_logger import SafeLogger
from api_connect_loop_detector import LoopDetector
from api_connect_dku_utils import get_dku_key_values, get_dku_duplicated_key_values, template_dict, format_template, is_reponse_xml, xml_to_json
from api_connect_dku_constants import DKUConstants
from api_connect_rest_api_auth import get_auth
from api_connect_retry_handler import DefaultRetryHandler


logger = SafeLogger("api-connect plugin", forbidden_keys=DKUConstants.FORBIDDEN_KEYS)
Expand All @@ -19,7 +20,7 @@ class RestAPIClientError(ValueError):

class RestAPIClient(object):

def __init__(self, credential, secure_credentials, endpoint, custom_key_values={}, session=None, behaviour_when_error=None):
def __init__(self, credential, secure_credentials, endpoint, custom_key_values={}, session=None, behaviour_when_error=None, retry_handler=None):
logger.info("Initialising RestAPIClient, credential={}, secure_credentials={}, endpoint={}".format(
logger.filter_secrets(credential),
logger.filter_secrets(secure_credentials),
Expand Down Expand Up @@ -134,6 +135,7 @@ def __init__(self, credential, secure_credentials, endpoint, custom_key_values={
self.secure_domain = "https://{}".format(self.secure_domain)
else:
self.session.auth = get_auth(credential)
self.retry_handler = retry_handler or DefaultRetryHandler()

def get(self, url, can_raise_exeption=True, **kwargs):
json_response = self.request("GET", url, can_raise_exeption=can_raise_exeption, **kwargs)
Expand All @@ -142,12 +144,10 @@ def get(self, url, can_raise_exeption=True, **kwargs):
def request(self, method, url, can_raise_exeption=True, **kwargs):
logger.info(u"Accessing endpoint {} with params={}".format(url, kwargs.get("params")))
self.assert_secure_domain(url)
self.enforce_throttling()
kwargs = template_dict(kwargs, **self.presets_variables)
if self.loop_detector.is_stuck_in_loop(url, kwargs.get("params", {}), kwargs.get("headers", {})):
raise RestAPIClientError("The api-connect plugin is stuck in a loop. Please check the pagination parameters.")
request_start_time = time.time()
self.time_last_request = request_start_time
error_message = None
status_code = None
response_headers = None
Expand Down Expand Up @@ -216,9 +216,17 @@ def request_with_cert(self, method, url, **kwargs):
)
tmp_key.seek(0)
kwargs["cert"] = (tmp_certificate.name, tmp_key.name)
response = self.session.request(method, url, **kwargs)
response = self.request_with_errors_retry(method, url, **kwargs)
return response
return self.session.request(method, url, **kwargs)
return self.request_with_errors_retry(method, url, **kwargs)

def request_with_errors_retry(self, method, url, **kwargs):
response = None
while self.retry_handler.should_retry(response):
self.enforce_throttling()
self.time_last_request = time.time()
response = self.session.request(method, url, **kwargs)
return response

def paginated_api_call(self, can_raise_exeption=True):
if self.pagination.params_must_be_blanked:
Expand Down Expand Up @@ -256,7 +264,7 @@ def start_paging(self):
self.pagination.reset_paging(counting_key=self.extraction_key, url=self.endpoint_url)

def enforce_throttling(self):
if self.time_between_requests and self.time_last_request:
if self.time_between_requests and self.time_last_request is not None:
current_time = time.time()
time_since_last_resquests = current_time - self.time_last_request
if time_since_last_resquests < self.time_between_requests:
Expand Down
Loading