Policy Guide
Data quality, reconciliation, and data cadence (freshness) policies in ADOC are catalog rules. This guide covers discovering policies, listing executions, starting runs (full, incremental, or selective), and reading status, results, and cancellation from Python.
Related product docs: For hard-linked vs. soft-linked policy behavior, see Pipeline Run Details in the ADOC User Guide documentation.
Get a Policy, List Policies, List Executions
from acceldata.client.adoc_client import AdocClientfrom acceldata.models.sdk.catalog import PolicyType, PolicyFilter, RuleTypeclient = AdocClient( url="https://<your-adoc-url>", access_key="<your-access-key>", secret_key="<your-secret-key>",)# Typed fetch (data quality, reconciliation, or data cadence)rule = client.get_policy(PolicyType.RECONCILIATION, "auth001_reconciliation")# Or untyped fetch by identifier → QualityRule (use rule_untyped.id for executions)rule_untyped = client.get_policy(identifier="my_rule_name")# List past executions by numeric rule id or by rule nameexecutions = client.policy_executions(1114, RuleType.DATA_QUALITY)executions = client.policy_executions("dq-scala", RuleType.DATA_QUALITY)# List executions for a rule already loaded: use the catalog rule id + RuleType, or a typed resourcerecon_rule = client.get_policy(PolicyType.RECONCILIATION, "auth001_reconciliation")recon_executions = client.policy_executions(recon_rule.rule.id, RuleType.RECONCILIATION)recon_res = client.get_recon_policy_resource("auth001_reconciliation")recon_executions_via_resource = recon_res.policy_executions(page=0, size=25)# List rules with a filterpolicy_filter = PolicyFilter(policy_type=RuleType.RECONCILIATION, enable=True)listed = client.list_all_policies(policy_filter=policy_filter, page=0, size=25)
Typed resources are also available for a single rule (data quality, reconciliation, or cadence):
dq_res = client.get_dq_policy_resource("my_dq")recon_res = client.get_recon_policy_resource("my_recon")cadence_res = client.get_cadence_policy_resource("my_cadence")# Example: list executions for that ruledq_res.policy_executions(page=0, size=25)
Starting a Run (Full, Incremental, or Selective)
Pass as the run payload on PolicyExecutionInput, execute_policy, execute_dq_rule, or execute_reconciliation_rule. Set execute_freshness_rule to choose how data is scoped:PolicyExecutionType
Type | Role |
| Scans the configured assets without per-run slice bounds (aside from optional Spark or rule-item options in the request). |
| Runs on data since the last persisted execution marker for that rule and its assets, as defined in the catalog. Many policies need only |
| Runs explicit slices. |
The sections below show , optional Spark settings, and related fields where they apply.markerConfigs
Selective runs always need one row per asset slice: the catalog asset ID, plus how to bound the data (numeric ID range, date/time window, file events, Kafka-style timestamps, and other marker types the catalog supports).
- For ID-bounded slices, use
.bounds_id_marker_config - For other cases, use the matching
helper in*_marker_configthat pairs with the generatedacceldata.services.marker_configclass (names align — for example,*MarkerConfig→BoundsDateTimeMarkerConfig), or setbounds_date_time_marker_configwithtype=(seemarker_type_for).API_TYPE_BY_MARKER_CLASS
Wrap the result in , and pass a MarkerConfig(...) as List[PolicyExecutionMarkerConfig] on the request.markerConfigs
execute_policy (Sync and Async) and the Run Handle
acceldata-sdk-python exposes , which executes policies both synchronously and asynchronously. The method returns an AdocClient.execute_policy instance, on which you can call Executor and get_result to obtain the execution's result and status.get_status
is always synchronous (it is not Python get_result()/async): the calling thread blocks in a poll loop until the run reaches a terminal status, until await non-terminal poll rounds are exhausted (default total_retries, unlimited), or until repeated result fetches fail (optional transient retries apply per fetch). -1 (default 5 seconds between polls) and sleep_interval apply to failure_strategy / Executor.get_result. get_execution_result only controls whether the client raises on terminal WARNING or ERROR outcomes; it does not make FailureStrategy.DoNotFail non-blocking.get_result
Parameters for Executing Policies (execute_policy)
The required parameters include:
Parameter | Description |
| A Boolean parameter that determines whether the policy runs synchronously or asynchronously. It is keyword-only and optional, with default |
| The policy type, specified as an enum parameter. Required. In acceldata-sdk-python, the parameter name is |
| The policy ID to execute, specified as a string or int parameter. Required (catalog rule ID or name, as the service accepts for that rule kind). |
| The run payload ( |
| In other Acceldata Python clients, this may appear as a Boolean specifying whether policy execution is incremental or full, often defaulting to |
| The run ID of the pipeline run where the policy is executing, specified as an optional int parameter (keyword-only). When set, it is sent as |
| An enum parameter that determines behavior when a failure occurs. Default: |
takes an enum of type failure_strategy, with three possible values:FailureStrategy
- DoNotFail: Terminal outcomes (including WARNING and ERROR result statuses) return an
without raising; HTTP or transport failures while polling yield a synthetic erroredExecutionResult. The same applies ifExecutionResultis exhausted while the run is still non-terminal.total_retries - FailOnError: Raises
when the result status indicates a failure other than WARNING. Also raisesPolicyExecutionCompletedWithErrorResultErroron HTTP or transport failures during polling when building that synthetic result. A completed run whose result status is WARNING does not raise underAcceldataSdkExceptionalone.FailOnError - FailOnWarning: Raises
when the result status is WARNING. It does not by itself raise on ERRORED or other error result statuses; usePolicyExecutionCompletedWithWarningErrorfor exceptions on failed (non-warning) completions.FailOnError
The execution result is available through on get_policy_execution_result, or AdocClient on the get_result instance returned from Executor; both return an execute_policy (or equivalent result object) when the run finishes.ExecutionResult
Additional keyword-only parameters on (not in the table above, but supported by this SDK):execute_policy
Parameter | Description |
| Seconds between status polls when the client waits inside |
| Cap on non-terminal poll rounds when waiting. Default |
| Optional |
Parameters to Retrieve Policy Execution Results
The required parameters include:
Parameter | Description |
| The policy type, specified as an enum parameter. Required. Accepts constant values such as |
| The execution ID to query for the result, specified as a string or int parameter. Required. |
| An enum parameter that determines behavior when a failure occurs. Default: |
Note: For more information on hard-linked and soft-linked policies, see the ADOC product documentation for rules and executions (also noted under Related product docs at the top of this guide).
also accepts an optional get_policy_execution_result for transient HTTP failures during polling. It waits until the run completes or the fetch fails — there is no poll-round cap on this method.transient_retry
Status: Client Methods and Executor.get_status
To get the current status, call on get_policy_execution_status, or AdocClient on the get_status returned from Executor. Both report status for the same execution ID as result polling.execute_policy
The required parameters for include:get_policy_execution_status
Parameter | Description |
| The policy type, specified as an enum parameter. Required. Accepts constant values such as |
| The execution ID to query for status, specified as a string or int parameter. Required. |
Note: on the get_status does not take any parameters; it uses the execution ID from the run that Executor (or execute_policy) already started.Executor.execute
Asynchronous Execution Example
The snippet assumes a configured client, as shown earlier in this guide.AdocClient
from acceldata.models.sdk.catalog import PolicyExecutionType, RuleTypefrom acceldata.models.sdk.catalog.policy_execution_request import PolicyExecutionInputfrom acceldata.models.sdk.catalog.executor import FailureStrategypolicy_execution_request = PolicyExecutionInput(PolicyExecutionType.FULL)async_executor = client.execute_policy( RuleType.DATA_QUALITY, 46, policy_execution_request, sync=False, failure_strategy=FailureStrategy.DoNotFail,)# Wait for execution to get the final resultexecution_result = async_executor.get_result(failure_strategy=FailureStrategy.DoNotFail)# Get the current statusexecution_status = async_executor.get_status()
Synchronous Execution Example
Same client assumption as the asynchronous example.
from acceldata.models.sdk.catalog import PolicyType, PolicyExecutionTypefrom acceldata.models.sdk.catalog.policy_execution_request import PolicyExecutionInputfrom acceldata.models.sdk.catalog.executor import FailureStrategypolicy_execution_request = PolicyExecutionInput(PolicyExecutionType.FULL)# This waits for execution to get the final result (within execute_policy when sync=True)sync_executor = client.execute_policy( RuleType.DATA_QUALITY, 46, policy_execution_request, sync=True, failure_strategy=FailureStrategy.DoNotFail,)# Wait for execution to get the final resultexecution_result = sync_executor.get_result(failure_strategy=FailureStrategy.DoNotFail)# Get the current statusexecution_status = sync_executor.get_status()
cancellation = sync_executor.cancel()
Cancel Execution Example
When the execution ID is already known (for example, from another system), use and get_policy_execution_status on get_policy_execution_result with that ID and the correct AdocClient.RuleType
Method | Parameters |
|
|
| Same, plus optional |
Stopping a run: call on the same cancel() handle. Typed policy resources may also expose a Executor helper for the same purpose.cancel_execution
Other ways to start a run: , execute_dq_rule, and execute_reconciliation_rule accept the same request shape; use the matching execute_freshness_rule or get_*_rule_result helpers to poll by execution ID.get_policy_execution_result
Trigger Policies: Execution Request
Note: Execution options (markers, Spark, rule items, and so on) depend on the ADOC catalog version. Keep your acceldata-sdk-python version aligned with the server.
Pass a to PolicyExecutionInput, execute_policy, execute_dq_rule, and execute_reconciliation_rule (the same parameter names appear in the code examples).execute_freshness_rule
Use to:PolicyExecutionInput
- Set how much data runs:
,FULL, orINCREMENTAL(throughSELECTIVE).PolicyExecutionType - For
, passSELECTIVEso each asset slice has bounds (see the ID, date/time, file event, and timestamp-based examples in this guide).markerConfigs - Optionally narrow which rule items run, link a pipeline run ID, tune Spark, or supply data quality dynamic SQL filter mappings.
PolicyExecutionInput Fields
All fields are keyword arguments except the first.
- executionType: Required first argument:
,PolicyExecutionType.FULL, orINCREMENTAL.SELECTIVE - markerConfigs: Optional. List of
; required forPolicyExecutionMarkerConfigin normal setups.SELECTIVE - ruleItemSelections: Optional. Int IDs of rule items to execute; omit to run all items. In this guide, runnable examples that set this field use
(see the ID-based data quality and reconciliation examples below).PolicyExecutionType.SELECTIVE - includeInQualityScore: Optional, default
. Includes this run in the quality score.True - pipelineRunId: Optional. Pipeline run ID to attach.
- sparkSQLDynamicFilterVariableMapping: Optional. For data quality policies with dynamic SQL filters.
- sparkResourceConfig: Optional. Spark resource limits and related settings.
- columnVariables: Optional. Column variable bindings when the policy requires them.
- sparkFilterSelectedColumns: Optional. Column names for Spark filter selection when applicable.
execute_dq_rule and execute_reconciliation_rule (Trigger Runs)
Use the same and PolicyExecutionType for data quality and reconciliation policies; only the client method changes (PolicyExecutionInput vs. execute_dq_rule).execute_reconciliation_rule
For data cadence (freshness), use with the same request shape.execute_freshness_rule
Loading a policy with returns a typed get_policy(PolicyType.*, identifier), DataQualityRuleResponse, or ReconciliationRuleResponse, each with a nested DataCadenceRuleResponse object — use rule in policy.rule.id (as shown in the snippets below). With execute_*_rule and no get_policy(identifier=...), the client returns a PolicyType; use QualityRule (there is no policy.id on that type).policy.rule
Data Quality Policy Execution Examples
Trigger full data quality policy: To run a full data quality policy across the entire dataset, use on PolicyExecutionType.FULL and call PolicyExecutionInput with the rule ID.execute_dq_rule
from acceldata.models.sdk.catalog import PolicyType, PolicyExecutionTypefrom acceldata.models.sdk.catalog.policy_execution_request import PolicyExecutionInputdq_policy = client.get_policy(PolicyType.DATA_QUALITY, "my_dq")req = PolicyExecutionInput(PolicyExecutionType.FULL)client.execute_dq_rule(dq_policy.rule.id, req)
Trigger incremental data quality policy: To run an incremental data quality policy based on a configured incremental strategy, use . The run advances from the catalog's last stored marker for that rule. The minimal request is often only the execution type; add PolicyExecutionType.INCREMENTAL only if the catalog requires an explicit marker for this policy.markerConfigs
dq_policy = client.get_policy(PolicyType.DATA_QUALITY, "my_dq")req = PolicyExecutionInput(PolicyExecutionType.INCREMENTAL)client.execute_dq_rule(dq_policy.rule.id, req)
Trigger selective data quality policy: To run a selective data quality policy over a subset of data, as determined by the chosen incremental strategy, and sparkSQLDynamicFilterVariableMapping are available, as shown in the optional example below. sparkResourceConfig applies only when the data quality policy includes SQL filters.sparkSQLDynamicFilterVariableMapping
ID-Based Selective Policy Execution
Uses a monotonically increasing column value to define data boundaries for policy execution, implemented with . Optional Spark tuning and BoundsIdMarkerConfig apply when the data quality policy uses dynamic SQL filters. Optionally pass sparkSQLDynamicFilterVariableMapping with catalog rule item definition IDs to run only those checks within the selective slice; omit the field to run every item on the policy. The same ruleItemSelections argument applies on ruleItemSelections and execute_reconciliation_rule objects.execute_freshness_rule[INLINE_CODE]PolicyExecutionInput[/INLINE_CODE]
from acceldata.models.sdk.catalog import PolicyType, PolicyExecutionTypefrom acceldata.models.sdk.catalog.policy_execution_request import PolicyExecutionInputfrom acceldata.models.api.catalog.policy_execution_marker_config import PolicyExecutionMarkerConfigfrom acceldata.models.api.catalog.marker_config import MarkerConfigfrom acceldata.models.api.catalog.spark_resource_config import SparkResourceConfigfrom acceldata.models.api.catalog.yunikorn_spark_resource_config import YunikornSparkResourceConfigfrom acceldata.models.api.catalog.rule_spark_sql_dynamic_filter_variable_mapping import ( RuleSparkSQLDynamicFilterVariableMapping,)from acceldata.models.api.catalog.mapping import Mappingfrom acceldata.services.marker_config import bounds_id_marker_configdq_policy = client.get_policy(PolicyType.DATA_QUALITY, "spark_sql_policy")bounds = bounds_id_marker_config( id_column_name="ID", from_id=0, to_id=1000,)marker_configs = [ # asset_id in the marker configuration refers to the unique identifier of the # underlying asset on which the Data Quality (DQ) policy is established. PolicyExecutionMarkerConfig( asset_id=9667404, marker_config=MarkerConfig(bounds), ),]yunikorn = YunikornSparkResourceConfig( min_executors=1, max_executors=2, executor_cores=2, executor_memory="2g", driver_cores=2, driver_memory="2g",)spark_cfg = SparkResourceConfig( invalidFields=[], yunikorn=yunikorn, additionalConfiguration={},)sql_maps = [ RuleSparkSQLDynamicFilterVariableMapping( rule_name="SelectiveDQPolicysparkSQLDynamicFilterVariable", mapping=[Mapping(key="column_name", is_column_variable=True, value="100")], ),]req = PolicyExecutionInput( PolicyExecutionType.SELECTIVE, markerConfigs=marker_configs, sparkResourceConfig=spark_cfg, sparkSQLDynamicFilterVariableMapping=sql_maps, ruleItemSelections=[101, 102], # replace with ids from the catalog / rule payload; omit for all items)client.execute_dq_rule(dq_policy.rule.id, req)
DateTime-Based Selective Policy Execution
Uses an increasing date column to define data boundaries for policy execution, implemented with .BoundsDateTimeMarkerConfig
from acceldata.models.sdk.catalog import PolicyType, PolicyExecutionTypefrom acceldata.models.sdk.catalog.policy_execution_request import PolicyExecutionInputfrom acceldata.models.api.catalog.policy_execution_marker_config import PolicyExecutionMarkerConfigfrom acceldata.models.api.catalog.marker_config import MarkerConfigfrom acceldata.services.marker_config import bounds_date_time_marker_configdq_policy = client.get_policy(PolicyType.DATA_QUALITY, "my_dq")bounds = bounds_date_time_marker_config( date_column_name="TO_DATE", format="yyyy-MM-dd", from_date="2023-07-01 00:00:00.000", to_date="2024-07-14 23:59:59.999", time_zone_id="Asia/Calcutta",)marker_configs = [ # asset_id in the marker configuration refers to the unique identifier of the # underlying asset on which the Data Quality (DQ) policy is established. PolicyExecutionMarkerConfig(asset_id=9667404, marker_config=MarkerConfig(bounds)),]req = PolicyExecutionInput(PolicyExecutionType.SELECTIVE, markerConfigs=marker_configs)client.execute_dq_rule(dq_policy.rule.id, req)
File Event-Based Selective Policy Execution
Uses file events to establish data boundaries for policy execution, implemented with . The BoundsFileEventMarkerConfig field must match a column on the bound asset (replace date_column_name with the actual column name).event_date
from acceldata.models.sdk.catalog import PolicyType, PolicyExecutionTypefrom acceldata.models.sdk.catalog.policy_execution_request import PolicyExecutionInputfrom acceldata.models.api.catalog.policy_execution_marker_config import PolicyExecutionMarkerConfigfrom acceldata.models.api.catalog.marker_config import MarkerConfigfrom acceldata.services.marker_config import bounds_file_event_marker_configdq_policy = client.get_policy(PolicyType.DATA_QUALITY, "my_dq")bounds = bounds_file_event_marker_config( date_column_name="event_date", from_date="2024-07-01 00:00:00.000", to_date="2024-07-01 23:59:59.999", time_zone_id="Asia/Calcutta",)marker_configs = [ # asset_id in the marker configuration refers to the unique identifier of the # underlying asset on which the Data Quality (DQ) policy is established. PolicyExecutionMarkerConfig(asset_id=1202688, marker_config=MarkerConfig(bounds)),]req = PolicyExecutionInput(PolicyExecutionType.SELECTIVE, markerConfigs=marker_configs)client.execute_dq_rule(dq_policy.rule.id, req)
Kafka Timestamp-Based Selective Policy Execution
Uses offsets associated with specified timestamps to set data boundaries for policy execution, implemented with .TimestampBasedMarkerConfig
from acceldata.models.sdk.catalog import PolicyType, PolicyExecutionTypefrom acceldata.models.sdk.catalog.policy_execution_request import PolicyExecutionInputfrom acceldata.models.api.catalog.policy_execution_marker_config import PolicyExecutionMarkerConfigfrom acceldata.models.api.catalog.marker_config import MarkerConfigfrom acceldata.services.marker_config import timestamp_based_marker_configdq_policy = client.get_policy(PolicyType.DATA_QUALITY, "kafka_dq_policy")bounds = timestamp_based_marker_config( format="yyyy-mm-dd", initial_offset="2023-06-01", time_zone_id="Asia/Calcutta",)marker_configs = [ # asset_id in the marker configuration refers to the unique identifier of the # underlying asset on which the Data Quality (DQ) policy is established. PolicyExecutionMarkerConfig(asset_id=5241961, marker_config=MarkerConfig(bounds)),]req = PolicyExecutionInput(PolicyExecutionType.SELECTIVE, markerConfigs=marker_configs)client.execute_dq_rule(dq_policy.rule.id, req)
Reconciliation Policy Execution Examples
Trigger full reconciliation policy: To trigger a full reconciliation policy across the entire dataset:
from acceldata.models.sdk.catalog import PolicyType, PolicyExecutionTypefrom acceldata.models.sdk.catalog.policy_execution_request import PolicyExecutionInputrecon_policy = client.get_policy(PolicyType.RECONCILIATION, "my_recon")req = PolicyExecutionInput(PolicyExecutionType.FULL)client.execute_reconciliation_rule(recon_policy.rule.id, req)
Trigger incremental reconciliation policy: To trigger an incremental reconciliation policy based on a configured incremental strategy:
recon_policy = client.get_policy(PolicyType.RECONCILIATION, "my_recon")req = PolicyExecutionInput(PolicyExecutionType.INCREMENTAL)client.execute_reconciliation_rule(recon_policy.rule.id, req)
Trigger selective reconciliation policy: To trigger a selective reconciliation policy over a subset of data, constrained by the selected incremental strategy:
ID-Based Selective Policy Execution
Uses a monotonically increasing column value to define data boundaries for policy execution, implemented with . The BoundsIdMarkerConfig in the marker configuration denotes the unique identifier of the underlying asset the reconciliation policy is based on. This can represent either the left or right asset ID.asset_id
from acceldata.models.sdk.catalog import PolicyType, PolicyExecutionTypefrom acceldata.models.sdk.catalog.policy_execution_request import PolicyExecutionInputfrom acceldata.models.api.catalog.policy_execution_marker_config import PolicyExecutionMarkerConfigfrom acceldata.models.api.catalog.marker_config import MarkerConfigfrom acceldata.models.api.catalog.spark_resource_config import SparkResourceConfigfrom acceldata.models.api.catalog.yunikorn_spark_resource_config import YunikornSparkResourceConfigfrom acceldata.services.marker_config import bounds_id_marker_configrecon_policy = client.get_policy(PolicyType.RECONCILIATION, "my_recon_policy")bounds = bounds_id_marker_config( id_column_name="ID", from_id=0, to_id=1000,)marker_configs = [ # asset_id in the marker configuration denotes the unique identifier of the # underlying asset on which the Reconciliation policy is based. This could # represent either the left or right asset ID. PolicyExecutionMarkerConfig(asset_id=9667404, marker_config=MarkerConfig(bounds)),]yunikorn = YunikornSparkResourceConfig( min_executors=1, max_executors=2, executor_cores=2, executor_memory="2g", driver_cores=2, driver_memory="2g",)spark_cfg = SparkResourceConfig( invalidFields=[], yunikorn=yunikorn, additionalConfiguration={},)req = PolicyExecutionInput( PolicyExecutionType.SELECTIVE, markerConfigs=marker_configs, sparkResourceConfig=spark_cfg, ruleItemSelections=[101, 102], # replace with ids from the catalog / rule payload; omit for all items)client.execute_reconciliation_rule(recon_policy.rule.id, req)
DateTime-Based Selective Policy Execution
Uses an increasing date column to define data boundaries for policy execution, implemented with .BoundsDateTimeMarkerConfig
from acceldata.models.sdk.catalog import PolicyType, PolicyExecutionTypefrom acceldata.models.sdk.catalog.policy_execution_request import PolicyExecutionInputfrom acceldata.models.api.catalog.policy_execution_marker_config import PolicyExecutionMarkerConfigfrom acceldata.models.api.catalog.marker_config import MarkerConfigfrom acceldata.services.marker_config import bounds_date_time_marker_configrecon_policy = client.get_policy(PolicyType.RECONCILIATION, "my_recon")bounds = bounds_date_time_marker_config( date_column_name="TO_DATE", format="yyyy-MM-dd", from_date="2023-07-01 00:00:00.000", to_date="2024-07-14 23:59:59.999", time_zone_id="Asia/Calcutta",)marker_configs = [ # asset_id in the marker configuration denotes the unique identifier of the # underlying asset on which the Reconciliation policy is based. This could # represent either the left or right asset ID. PolicyExecutionMarkerConfig(asset_id=9667404, marker_config=MarkerConfig(bounds)),]req = PolicyExecutionInput(PolicyExecutionType.SELECTIVE, markerConfigs=marker_configs)client.execute_reconciliation_rule(recon_policy.rule.id, req)
File Event-Based Selective Policy Execution
Uses file events to establish data boundaries for policy execution, implemented with . Use a real BoundsFileEventMarkerConfig for the backing asset (same note as for data quality file-event markers).date_column_name
from acceldata.models.sdk.catalog import PolicyType, PolicyExecutionTypefrom acceldata.models.sdk.catalog.policy_execution_request import PolicyExecutionInputfrom acceldata.models.api.catalog.policy_execution_marker_config import PolicyExecutionMarkerConfigfrom acceldata.models.api.catalog.marker_config import MarkerConfigfrom acceldata.services.marker_config import bounds_file_event_marker_configrecon_policy = client.get_policy(PolicyType.RECONCILIATION, "my_recon")bounds = bounds_file_event_marker_config( date_column_name="event_date", from_date="2024-07-01 00:00:00.000", to_date="2024-07-01 23:59:59.999", time_zone_id="Asia/Calcutta",)marker_configs = [ # asset_id in the marker configuration denotes the unique identifier of the # underlying asset on which the Reconciliation policy is based. This could # represent either the left or right asset ID. PolicyExecutionMarkerConfig(asset_id=1202688, marker_config=MarkerConfig(bounds)),]req = PolicyExecutionInput(PolicyExecutionType.SELECTIVE, markerConfigs=marker_configs)client.execute_reconciliation_rule(recon_policy.rule.id, req)
Kafka Timestamp-Based Selective Policy Execution
Uses offsets associated with specified timestamps to set data boundaries for policy execution, implemented with .TimestampBasedMarkerConfig
from acceldata.models.sdk.catalog import PolicyType, PolicyExecutionTypefrom acceldata.models.sdk.catalog.policy_execution_request import PolicyExecutionInputfrom acceldata.models.api.catalog.policy_execution_marker_config import PolicyExecutionMarkerConfigfrom acceldata.models.api.catalog.marker_config import MarkerConfigfrom acceldata.services.marker_config import timestamp_based_marker_configrecon_policy = client.get_policy(PolicyType.RECONCILIATION, "kafka_recon_policy")bounds = timestamp_based_marker_config( format="yyyy-mm-dd", initial_offset="2023-06-01", time_zone_id="Asia/Calcutta",)marker_configs = [ # asset_id in the marker configuration denotes the unique identifier of the # underlying asset on which the Reconciliation policy is based. This could # represent either the left or right asset ID. PolicyExecutionMarkerConfig(asset_id=5241961, marker_config=MarkerConfig(bounds)),]req = PolicyExecutionInput(PolicyExecutionType.SELECTIVE, markerConfigs=marker_configs)client.execute_reconciliation_rule(recon_policy.rule.id, req)
What's Next
After you complete this section, explore:
- Tags and Labels – Learn how to attach tags and labels to the assets these policies run against.
- Pipelines Guide – Learn how to link a policy execution to a pipeline run using
.pipelineRunId

Have a suggestion?