Python API
Accuracy
from ai.chronon.types import Accuracy Aggregation
from ai.chronon.types import Aggregation
Aggregation(
input_column: str = None,
operation: Union[Operation, OperationWithArgs] = None,
windows: Union[List[Window], List[str]] = None,
buckets: List[str] = None,
tags: Dict[str, str] = None,
filter: str = None,
column_alias: str = None
) -> ttypes.AggregationParameters:
input_column(str) — Column on which the aggregation needs to be performed. This should be one of the input columns specified on the keys of theselectin theQuery'sSourceoperation(ttypes.Operation) — Operation to use to aggregate the input columns. For example, MAX, MIN, COUNT Some operations have arguments, like last_k, approx_percentiles etc., Defaults to "LAST".windows(List[common.Window]) — Length to window to calculate the aggregates on. Strings like "1h", "30d" are also accepted. Minimum window size is 1hr. Maximum can be arbitrary. When not defined, the computation is un-windowed.buckets(List[str]) — Besides the GroupBy.keys, this is another level of keys for use under this aggregation. Using this would create an output as a map of string to aggregate.filter(str) — Optional SQL boolean expression restricting which rows feed this aggregation. Compiled away in the Python layer into aCASE WHEN (<filter>) THEN <input_column> ELSE NULL ENDcolumn on the sourceselects(non-matching rows become NULL and are skipped by every operation). Use AND/OR/parentheses for compound conditions. Requirescolumn_aliasso the output feature has an explicit, collision-free name.column_alias(str) — Output-name prefix, replacinginput_column; the feature becomes<column_alias>_<op>_<window>. Only valid together withfilter.
Returns: An aggregate defined with the specified operation.
Aggregation
from ai.chronon.types import Aggregation Aggregations
from ai.chronon.types import Aggregations
Aggregations(**agg_dict) BootstrapPart
from ai.chronon.types import BootstrapPart
BootstrapPart(
table: str,
key_columns: List[str] = None,
query: Query = None
) -> ttypes.BootstrapPartBootstrap is the concept of using pre-computed feature values and skipping backfill computation during the training data generation phase. Bootstrap can be used for many purposes:
- Generating ongoing feature values from logs
- Backfilling feature values for external features (in which case Chronon is unable to run backfill)
- Initializing a new Join by migrating old data from an older Join and reusing data
One can bootstrap against any of these:
- join part fields: Bootstrap can happen at individual field level within a join part. If all fields within a group by are bootstrapped, then we skip computation for group by. Otherwise, the whole thing will be re-run but only the values for the non-bootstrapped fields will be retained in the final table.
- external part fields: Bootstrap can happen at individual field level within an external part. Since there is no backfill logic in chronon for external part, all non-bootstrapped fields in external parts are left as NULLs.
- derivation fields: Derived fields can also be bootstrapped. Since derived fields depend on "base" fields (either join part or external part), chronon will try to trigger the least amount of computation possible. For example, if there is a join part where all derived fields that depend on the join part have been bootstrapped, then we skip the computation for this join part.
- keys: Keys of both join parts and external parts can be bootstrapped. During offline table generation, we will first try to utilize key's data from left table; if it's not there, then we utilize bootstrap. For contextual features, we also support propagating the key bootstrap to the values.
Parameters:
table— Name of hive table that contains feature values where rows are 1:1 mapped to left tablekey_columns— Keys to join bootstrap table to left tablequery— Selected columns (features & keys) and filtering conditions of the bootstrap tables.
CaseWhen
from ai.chronon.types import CaseWhenFluent builder for a SQL CASE WHEN expression.
Renders to a plain SQL string, so it can be used directly as a selects value (or anywhere a SQL expression string is expected):
when("status = 'failed'", "amount").when("status = 'pending'", "0").else_(-1)
# -> "CASE WHEN status = 'failed' THEN amount WHEN status = 'pending' THEN 0 ELSE -1 END"
Omitting else_ leaves the default as SQL NULL. condition and value are raw SQL expressions (a Python None renders as NULL);
string literals must be quoted by the caller (e.g. "'SMS'").
| Name | Description |
|---|---|
obj | |
parts | |
otherwise |
ContextualSource
from ai.chronon.types import ContextualSource
ContextualSource(fields: FieldsType, team="default") -> ttypes.ExternalSourceContextual source values are passed along for logging. No external request is actually made.
DefaultAggregation
from ai.chronon.types import DefaultAggregation
DefaultAggregation(keys, sources, operation=Operation.LAST, tags=None) DeploymentSpec
from ai.chronon.types import DeploymentSpec| Name | Description |
|---|---|
containerConfig | |
endpointConfig | |
resourceConfig | |
rolloutStrategy |
DeploymentStrategyType
from ai.chronon.types import DeploymentStrategyType| Name | Description |
|---|---|
BLUE_GREEN | |
ROLLING | |
IMMEDIATE |
Derivation
from ai.chronon.types import Derivation
Derivation(name: str, expression: str) -> ttypes.DerivationDerivation allows arbitrary SQL select clauses to be computed using columns from the output of a GroupBy or Join. The results will be available both in online fetching response maps and in offline Hive tables.
For Joins, column names are automatically constructed according to the convention: {join_part_prefix}_{group_by_name}_{input_column_name}_{aggregation_operation}_{window}_{by_bucket} (prefix, window, and bucket are optional).
For ExternalParts, column names follow: ext_{external_source_name}_{value_column}.
Note that only values can be used in derivations, not keys. If you want to use a key in the
derivation, you must define it as a contextual field and refer to it with its prefix included,
for example: ext_contextual_request_id.
If both name and expression are set to "*", then every raw column will be included along with the derived columns.
Parameters:
name— output column name of the SQL expressionexpression— any valid Spark SQL select clause based on joinPart or externalPart columns
Returns: a Derivation object representing a single derived column or a wildcard ("*") selection.
Derivation
from ai.chronon.types import Derivation
Derivation(name: str, expression: str) -> ttypes.Derivation.. deprecated:
Use ``from ai.chronon.types import Derivation`` instead. Derivation
from ai.chronon.types import Derivation
Derivation(name: str, expression: str) -> ttypes.Derivation.. deprecated:
Use ``from ai.chronon.types import Derivation`` instead. EndpointConfig
from ai.chronon.types import EndpointConfig| Name | Description |
|---|---|
endpointName | |
additionalConfigs |
EngineType
from ai.chronon.types import EngineType| Name | Description |
|---|---|
SPARK | |
BIGQUERY | |
SNOWFLAKE |
EntitySource
from ai.chronon.types import EntitySource
EntitySource(
snapshot_table: str,
query: Query,
mutation_table: str = None,
mutation_topic: str = None
) -> ttypes.SourceEntity Sources represent data that gets mutated over-time - at row-level. This is a group of three data elements.
snapshotTable, mutationTable and mutationTopic. mutationTable and mutationTopic are only necessary if we are trying
to create realtime or point-in-time aggregations over these sources. Entity sources usually map 1:1 with a database
tables in your OLTP store that typically serves live application traffic. When mutation data is absent they map 1:1
to dim tables in star schema.
Parameters:
snapshotTable— Snapshot table currently needs to be a 'ds' (date string - yyyy-MM-dd) partitioned hive table.mutationTable— Topic is a kafka table. The table contains all the events that historically came through this topic. We need all the fields present in the snapshot table, PLUS two additional fields,mutation_time- milliseconds since epoch of type Long that represents the time of the mutationis_before- a boolean flag that represents whether this row contains values before or after the mutation.mutationTopic— The logic used to scan both the table and the topic. Contains row level transformations and filtering expressed as Spark SQL statements.query— If each new hive partition contains not just the current day's events but the entire set of events since the begininng. The key property is that the events are not mutated across partitions.
EventSource
from ai.chronon.types import EventSource
EventSource(
table: str,
query: Query,
topic: str = None,
is_cumulative: bool = None
) -> ttypes.SourceEvent Sources represent data that gets generated over-time. Typically, but not necessarily, logged to message buses like kafka, kinesis or google pub/sub. fct tables are also event source worthy.
Parameters:
table— Table currently needs to be a 'ds' (date string - yyyy-MM-dd) partitioned hive table. Table names can contain subpartition specs, example db.table/system=mobile/currency=USDtopic— Topic is a kafka table. The table contains all the events historically came through this topic.query— The logic used to scan both the table and the topic. Contains row level transformations and filtering expressed as Spark SQL statements.isCumulative— If each new hive partition contains not just the current day's events but the entire set of events since the begininng. The key property is that the events are not mutated across partitions.
ExternalPart
from ai.chronon.types import ExternalPart
ExternalPart(
source: ExternalSource,
key_mapping: Dict[str, str] = None,
prefix: str = None
) -> ttypes.ExternalPartUsed to describe which ExternalSources to pull features from while fetching online. This data also goes into logs based on sample percent.
Just as in JoinPart, key_mapping is used to map the join left's keys to
external source's keys. "vendor" and "buyer" on left side (query map)
could both map to a "user" in an account data external source. You would
create one ExternalPart for vendor with params: (key_mapping={vendor: user}, prefix=vendor) another for buyer.
This doesn't have any implications offline besides logging. "right_parts" can be both backfilled and logged. Whereas, "external_parts" can only be logged. If you need the ability to backfill an external source, look into creating an EntitySource with mutation data for point-in-time-correctness.
Parameters:
source— External source to join withkey_mapping— How to map the keys from the query/left side to the sourceprefix— Sometime you want to use the same source to fetch data for different entities in the query. Eg., A transaction between a buyer and a seller might query "user information" serivce/source that has information about both buyer & seller
ExternalSource
from ai.chronon.types import ExternalSource
ExternalSource(
name: str,
team: str,
key_fields: FieldsType,
value_fields: FieldsType
) -> ttypes.ExternalSourceExternal sources are online only data sources. During fetching, using chronon java client, they consume a Request containing a key map (name string to value). And produce a Response containing a value map.
This is primarily used in Joins. We also expose a fetchExternal method in java client library that can be used to fetch a batch of External source requests efficiently.
Internally Chronon will batch these requests to the service and parallelize fetching from different services, while de-duplicating given a batch of join requests.
The implementation of how to fetch is an ExternalSourceHandler in
scala/java api that needs to be registered while implementing
ai.chronon.online.Api with the name used in the ExternalSource. This is
meant for re-usability of external source definitions.
Parameters:
name— name of the external source to fetch from. Should match the name in the registry.key_fields— List of tuples of string and DataType. This is what will be given to ExternalSource handler registered in Java API. Eg.,[('key1', DataType.INT, 'key2', DataType.STRING)]value_fields— List of tuples of string and DataType. This is what the ExternalSource handler will respond with:[ ('value0', DataType.INT), ('value1', DataType.MAP(DataType.STRING, DataType.LONG), ('value2', DataType.STRUCT( name = 'Context', ('field1', DataType.INT), ('field2', DataType.DOUBLE) )) ]
GroupBy
from ai.chronon.types import GroupBy
GroupBy(
sources: Union[Sequence[Source], Source],
keys: List[str],
aggregations: Optional[List[Aggregation]],
version: Optional[int] = None,
derivations: List[Derivation] = None,
accuracy: Accuracy = None,
output_namespace: str = None,
table_properties: Dict[str, str] = None,
tags: Dict[str, str] = None,
online: bool = None,
online_strategy: Optional[OnlineStrategy] = None,
production: bool = None,
key_filter: Optional[Union[Source, EntitySource]] = None,
offline_schedule: str = None,
online_schedule: Optional[str] = None,
conf: ConfigProperties = None,
env_vars: EnvironmentVariables = None,
cluster_conf: ClusterConfigProperties = None,
step_days: int = None,
disable_historical_backfill: bool = False,
environments: Optional[List[str]] = None,
partition_interval: Optional[Union[Window, str]] = None,
partition_offset: Optional[Union[Window, str]] = None,
workflow_concurrency: Optional[int] = None
) -> ttypes.GroupByParameters:
version— TODOsources(List of sources or a single source) — can be constructed as entities or events or joinSource:import gen_thrift.api.ttypes as chronon events = chronon.Source(events=chronon.Events( table=YOUR_TABLE, topic=YOUR_TOPIC # <- OPTIONAL for serving query=chronon.Query(...) isCumulative=False # <- defaults to false. )) Or entities = chronon.Source(entities=chronon.Entities( snapshotTable=YOUR_TABLE, mutationTopic=YOUR_TOPIC, mutationTable=YOUR_MUTATION_TABLE query=chronon.Query(...) )) or joinSource = chronon.Source(joinSource=chronon.JoinSource( join = YOUR_CHRONON_PARENT_JOIN, query = chronon.Query(...) ))Multiple sources can be supplied to backfill the historical values with their respective start and end partitions. However, only one source is allowed to be a streaming one.keys(List[String]) — List of primary keys that defines the data that needs to be collected in the result table. Similar to the GroupBy in the SQL context.aggregations(List[gen_thrift.api.ttypes.Aggregation]) — List of aggregations that needs to be computed for the data following the grouping defined by the keys:import gen_thrift.api.ttypes as chronon aggregations = [ chronon.Aggregation(input_column="entity", operation=Operation.LAST), chronon.Aggregation(input_column="entity", operation=Operation.LAST, windows=['7d']) ],online(bool) — Should we upload the result data of this conf into the KV store so that we can fetch/serve this GroupBy online. Once Online is set to True, you ideally should not change the conf.production(bool) — This when set can be integrated to trigger alerts. You will have to integrate this flag into your alerting system yourself.env(Dict[str, Dict[str, str]]) — This is a dictionary of "mode name" to dictionary of "env var name" to "env var value":{ 'backfill' : { 'VAR1' : 'VAL1', 'VAR2' : 'VAL2' }, 'upload' : { 'VAR1' : 'VAL1', 'VAR2' : 'VAL2' } 'streaming' : { 'VAR1' : 'VAL1', 'VAR2' : 'VAL2' } }These vars then flow into run.py and the underlying spark_submit.sh. These vars can be set in other places as well. The priority order (descending) is as below 1. env vars set while using run.py "VAR=VAL run.py --mode=backfill <name>" 2. env vars set here in Join's env param 3. env vars set in `team.json['team.production.<MODE NAME>']` 4. env vars set in `team.json['default.production.<MODE NAME>']`table_properties(Dict[str, str]) — Specifies the properties on output hive tables. Can be specified in teams.json.output_namespace(str) — In backfill mode, we will produce data into hive. This represents the hive namespace that the data will be written into. You can set this at the teams.json level.accuracy(gen_thrift.api.ttypes.SNAPSHOT or gen_thrift.api.ttypes.TEMPORAL) — Defines the computing accuracy of the GroupBy. If "Snapshot" is selected, the aggregations are computed based on the partition identifier - "ds" time column. If "Temporal" is selected, the aggregations are computed based on the event time - "ts" time column.lag(int) — Param that goes into customJson. You can pull this out of the json at path "metaData.customJson.lag" This is used by airflow integration to pick an older hive partition to wait on.offline_schedule(str) — The offline schedule interval for batch jobs. Supports standard cron expressions, including regular sub-daily schedules. Examples:'@daily': Legacy format for midnight daily execution '0 2 * * *': Daily at 2:00 AM '0 */3 * * *': Every 3 hours '5/15 * * * *': Every 15 minutes with a 5 minute processing offset '@never': Explicitly disable offline schedulingonline_schedule(Optional[str]) — The online schedule interval for real-time serving jobs. Supports standard cron expressions When online=True and online_schedule is not specified, defaults to offline_schedule when present, otherwise "@daily". Set to "@never" to explicitly disable online scheduling even when online=True. Examples follow the same format as offline_schedule.partition_interval(Optional[Union[common.Window, str]]) — Output partition interval for this GroupBy. Examples: "1d", "3h", "15m". When set below daily, Chronon uses "yyyy-MM-dd-HH-mm" ds values.partition_offset(Optional[Union[common.Window, str]]) — Offset from UTC midnight/epoch for the output partition grid. Defaults to zero (boundaries at midnight) and must be declared explicitly to move the grid; the cron fire phase is treated as a processing delay relative to the declared grid, never as a grid offset.tags(Dict[str, str]) — Additional metadata that does not directly affect feature computation, but is useful to track for management purposes.derivations(List[gen_thrift.api.ttypes.Drivation]) — Derivation allows arbitrary SQL select clauses to be computed using columns from the output of group by backfill output schema. It is supported for offline computations for now.kwargs(Dict[str, str]) — Additional properties that would be passed to run.py if specified under additional_args property. And provides an option to pass custom values to the processing logic.conf— Configuration properties for the GroupBy. Depending on the mode we layer confs with the following priority: 1. conf set in the GroupBy.conf.2. conf set in the GroupBy.conf.common 3. conf set in the team.conf. 4. conf set in the team.conf.common 5. conf set in the default.conf. 6. conf set in the default.conf.common env_vars— Environment variables for the GroupBy. Depending on the mode we layer envs with the following priority: 1. env vars set in the GroupBy.env.2. env vars set in the GroupBy.env.common 3. env vars set in the team.env. 4. env vars set in the team.env.common 5. env vars set in the default.env. 6. env vars set in the default.env.common cluster_conf— Cluster configuration properties for the join.step_days— The maximum number of days to output at oncekey_filter(gen_thrift.api.ttypes.EntitySource (or a Source wrapping one)) — An entities source whose snapshot partition for the upload date restricts which keys make it into the batch upload: the aggregated upload rows (one per key) are semi-joined against the distinct key tuples found in this source. The filter's query.selects must produce columns named after (a subset of) the GroupBy keys. Only applied by uploads to shrink upload size - backfills ignore it, and it does not affect semantic hashes, so setting or changing it never re-triggers jobs. Note: for TEMPORAL GroupBys with a streaming topic, the streaming job still writes all keys, so filtered-out keys may serve partial streaming-only aggregates instead of nulls.environments(List[str]) — List of environments where this GroupBy should be deployed/available. Defaults to ['prod']. Valid values: 'prod', 'canary' (case-insensitive).workflow_concurrency— Default maximum number of workflow steps Hub may allocate concurrently when a workflow is started from this GroupBy. Request-level overrides take precedence.
Returns: A GroupBy object containing specified aggregations.
GroupBy
from ai.chronon.types import GroupBy
GroupBy(
sources: Union[Sequence[Source], Source],
keys: List[str],
aggregations: List[Aggregation],
row_id_column: str,
valid_from_column: str,
valid_to_column: str,
version: Optional[int] = None,
accuracy: Accuracy = Accuracy.TEMPORAL,
output_namespace: str = None,
table_properties: Dict[str, str] = None,
tags: Dict[str, str] = None,
online: bool = group_by.None,
production: bool = group_by.None,
offline_schedule: str = None,
online_schedule: Optional[str] = None,
conf: ConfigProperties = None,
env_vars: EnvironmentVariables = None,
cluster_conf: ClusterConfigProperties = None,
step_days: int = None,
disable_historical_backfill: bool = False,
as_of_column: str = "ts"
) -> ttypes.GroupBy InferenceSpec
from ai.chronon.types import InferenceSpec| Name | Description |
|---|---|
modelBackend | |
modelBackendParams | |
resourceConfig |
Join
from ai.chronon.types import Join
Join(
left: Source,
right_parts: List[JoinPart],
row_ids: Optional[Union[str, List[str]]] = None,
version: Optional[int] = None,
online_external_parts: List[ExternalPart] = None,
bootstrap_parts: List[BootstrapPart] = None,
bootstrap_from_log: bool = False,
skew_keys: Dict[str, List[str]] = None,
derivations: List[Derivation] = None,
output_namespace: str = None,
table_properties: Dict[str, str] = None,
online: bool = False,
production: bool = False,
sample_percent: float = 100.0,
check_consistency: bool = None,
consistency_sample_percent: float = 5.0,
use_long_names: bool = False,
offline_schedule: str = "@daily",
online_schedule: str = None,
historical_backfill: bool = None,
conf: ConfigProperties = None,
env_vars: EnvironmentVariables = None,
cluster_conf: ClusterConfigProperties = None,
step_days: int = None,
enable_stats_compute: bool = None,
modular_execution: bool = False,
environments: Optional[List[str]] = None,
partition_interval: Optional[Union[Window, str]] = None,
partition_offset: Optional[Union[Window, str]] = None,
workflow_concurrency: Optional[int] = None
) -> ttypes.JoinConstruct a join object. A join can pull together data from various GroupBy's both offline and online. This is also the focal point for logging, data quality computation and monitoring. A join maps 1:1 to models in ML usage.
Parameters:
left(ai.chronon.api.Source) — The source on the left side, when Entities, all GroupBys are join with SNAPSHOT accuracy (midnight values). When left is events, if on the right, either when GroupBy's are TEMPORAL, or when topic is specified, we perform a TEMPORAL / point-in-time join.right_parts(List[ai.chronon.api.JoinPart]) — The list of groupBy's to join with. GroupBy's are wrapped in a JoinPart, which contains additional information on how to join the left side with the GroupBy.check_consistency(bool) — If online serving data should be compared with backfill data - as online-offline-consistency metrics. The metrics go into hive and your configured kv store for further visualization and monitoring.additional_args(List[str]) — Additional args go intocustomJsonofai.chronon.api.MetaDatawithin theai.chronon.api.Joinobject. This is a place for arbitrary information you want to tag your conf with.additional_env(List[str]) — Deprecated, see envonline(bool) — Should we upload this conf into kv store so that we can fetch/serve this join online. Once Online is set to True, you ideally should not change the conf.production(bool) — This when set can be integrated to trigger alerts. You will have to integrate this flag into your alerting system yourself.output_namespace(str) — In backfill mode, we will produce data into hive. This represents the hive namespace that the data will be written into. You can set this at the teams.json level.table_properties— Specifies the properties on output hive tables. Can be specified in teams.json.lag— Param that goes into customJson. You can pull this out of the json at path "metaData.customJson.lag" This is used by airflow integration to pick an older hive partition to wait on.skew_keys— While back-filling, if there are known irrelevant keys - like user_id = 0 / NULL etc. You can specify them here. This is used to blacklist crawlers etcsample_percent— Online only parameter. What percent of online serving requests to this join should be logged into warehouse.consistency_sample_percent— Online only parameter. What percent of online serving requests to this join should be sampled to compute online offline consistency metrics. if sample_percent=50.0 and consistency_sample_percent=10.0, then basically the consistency job runs on 5% of total traffic.online_external_parts— users can register external sources into Api implementation. Chronon fetcher can invoke the implementation. This is applicable only for online fetching. Offline this will not be produce any values.offline_schedule— Schedule expression for offline join compute tasks. Supports standard cron expressions, including regular sub-daily schedules. Examples: '@daily' (midnight), '0 2 * * *' (2am daily), '0 */3 * * *' (every 3 hours), '5/15 * * * *' (every 15 minutes with a 5 minute processing offset). Use '@never' to explicitly disable offline scheduling.online_schedule— Schedule expression for online/deploy tasks. When online=True and online_schedule is not specified, defaults to offline_schedule. Set to "@never" to explicitly disable online scheduling even when online=True. Supports the same format as offline_schedule.partition_interval— Output partition interval for this Join. Examples: "1d", "3h", "15m". When set below daily, Chronon uses "yyyy-MM-dd-HH-mm" ds values.partition_offset— Offset from UTC midnight/epoch for the output partition grid. Defaults to zero (boundaries at midnight) and must be declared explicitly to move the grid; the cron fire phase is treated as a processing delay relative to the declared grid, never as a grid offset.row_ids— Columns of the left table that uniquely define a training record. Used as default keys during bootstrap. Optional.bootstrap_parts— A list of BootstrapPart used for the Join. See BootstrapPart doc for more detailsbootstrap_from_log— If set to True, will use logging table to generate training data by default and skip continuous backfill. Logging will be treated as another bootstrap source, but other bootstrap_parts will take precedence.historical_backfill(bool) — Flag to indicate whether join backfill should backfill previous holes. Setting to false will only backfill latest single partitionconf— Configuration properties for the join. Depending on the mode we layer confs with the following priority: 1. conf set in the join.conf.2. conf set in the join.conf.common 3. conf set in the team.conf. 4. conf set in the team.conf.common 5. conf set in the default.conf. 6. conf set in the default.conf.common env_vars— Environment variables for the join. Depending on the mode we layer envs with the following priority: 1. env vars set in the join.env.2. env vars set in the join.env.common 3. env vars set in the team.env. 4. env vars set in the team.env.common 5. env vars set in the default.env. 6. env vars set in the default.env.common cluster_conf— Cluster configuration properties for the join.step_days— The maximum number of days to output at onceenable_stats_compute(bool) — Whether to enable enhanced statistics computation and upload for this join. When True, stats compute and upload nodes will be added to the workflow.modular_execution(bool) — When True, uses modular join planning (JoinPlanner) instead of the default monolith planner (MonolithJoinPlanner).environments(List[str]) — List of environments where this join should be deployed/available. Defaults to ['prod']. Valid values: 'prod', 'canary' (case-insensitive).workflow_concurrency(int) — Default maximum number of workflow steps Hub may allocate concurrently when a workflow is started from this Join. Request-level overrides take precedence.
Returns: A join object that can be used to backfill or serve data. For ML use-cases this should map 1:1 to model.
JoinPart
from ai.chronon.types import JoinPart
JoinPart(
group_by: GroupBy,
key_mapping: Dict[str, str] = None,
prefix: str = None,
tags: Dict[str, str] = None
) -> ttypes.JoinPartSpecifies HOW to join the left of a Join with GroupBy's.
Parameters:
group_by(ai.chronon.api.GroupBy) — The GroupBy object to join with. Keys on left are used to equi join with keys on right. When left is entities all GroupBy's are computed as of midnight. When left is events, we do a point-in-time join when right.accuracy == TEMPORAL OR right.source.topic != nullkey_mapping(Dict[str, str]) — Names of keys don't always match on left and right, this mapping tells us how to map when they don't.prefix— All the output columns of the groupBy will be prefixed with this string. This is used when you need to join the same groupBy more than once withleft. Say on the left you have seller and buyer, on the group you have a user's avg_price, and you want to join the left (seller, buyer) with (seller_avg_price, buyer_avg_price) you would use key_mapping and prefix parameters.tags(Dict[str, str]) — Additional metadata about the JoinPart that you wish to track. Does not effect computation.
Returns: JoinPart specifies how the left side of a join, or the query in online setting, would join with the right side components like GroupBys.
JoinSource
from ai.chronon.types import JoinSource
JoinSource(join: Join, query: Query = None) -> ttypes.SourceThe output of a join can be used as a source for GroupBy.
Useful for expressing complex computation in chronon.
Offline this simply means that we will compute the necessary date ranges of the join
before we start computing the GroupBy.
Online we will:
- enrich the stream/topic of
join.leftwith all the columns defined by the join - apply the selects & wheres defined in the
query - perform aggregations defined in the downstream
GroupBy - write the result to the kv store.
Metric
from ai.chronon.types import Metric| Name | Description |
|---|---|
name | |
threshold |
Model
from ai.chronon.types import Model
Model(
version: str,
inference_spec: Optional[InferenceSpec] = None,
input_mapping: Optional[Dict[str, str]] = None,
output_mapping: Optional[Dict[str, str]] = None,
value_fields: Optional[FieldsType] = None,
model_artifact_base_uri: Optional[str] = None,
training_conf: Optional[TrainingSpec] = None,
deployment_conf: Optional[DeploymentSpec] = None,
output_namespace: Optional[str] = None,
table_properties: Optional[Dict[str, str]] = None,
tags: Optional[Dict[str, str]] = None,
environments: Optional[List[str]] = None,
partition_interval: Optional[Union[Window, str]] = None,
partition_offset: Optional[Union[Window, str]] = None,
workflow_concurrency: Optional[int] = None
) -> ttypes.ModelCreates a Model object for ML model inference and orchestration.
Parameters:
version(str) — Version string for the model configurationinference_spec(InferenceSpec) — Model + model backend specific details necessary to perform inferenceinput_mapping(Dict[str, str]) — Spark SQL queries to transform input data to the format expected by the modeloutput_mapping(Dict[str, str]) — Spark SQL queries to transform model output to desired output formatvalue_fields(FieldsType) — List of tuples of (field_name, DataType) defining the schema of the model's output values. If provided, creates a STRUCT schema that will be set as the model's valueSchema. Example: [('score', DataType.DOUBLE), ('category', DataType.STRING)]model_artifact_base_uri(str) — Base URI where trained model artifacts are storedtraining_conf(TrainingSpec) — Configs related to orchestrating model training jobsdeployment_conf(DeploymentSpec) — Configs related to orchestrating model deploymentoutput_namespace(str) — Namespace for the model outputtable_properties(Dict[str, str]) — Additional table properties for the model outputtags(Dict[str, str]) — Additional metadata that does not directly affect computation, but is useful for management.environments(List[str]) — List of environments where this Model should be deployed/available. Defaults to ['prod']. Valid values: 'prod', 'canary' (case-insensitive).partition_interval(Optional[Union[common.Window, str]]) — Output partition interval for model training/deploy nodes. Examples: "1d", "3h", "15m". When set below daily, Chronon uses "yyyy-MM-dd-HH-mm" ds values.partition_offset(Optional[Union[common.Window, str]]) — Offset from UTC midnight/epoch for the output partition grid.workflow_concurrency(int) — Default maximum number of workflow steps Hub may allocate concurrently when a workflow is started from this Model. Request-level overrides take precedence.
Returns: A Model object
ModelBackend
from ai.chronon.types import ModelBackend| Name | Description |
|---|---|
VERTEXAI | |
SAGEMAKER | |
RAY |
ModelTransforms
from ai.chronon.types import ModelTransforms
ModelTransforms(
sources: Sequence[ANY_SOURCE_TYPE],
models: List[Model],
version: int,
passthrough_fields: Optional[List[str]] = None,
key_fields: Optional[FieldsType] = None,
output_namespace: Optional[str] = None,
table_properties: Optional[Dict[str, str]] = None,
tags: Optional[Dict[str, str]] = None,
environments: Optional[List[str]] = None,
workflow_concurrency: Optional[int] = None
) -> ttypes.ModelTransformsModelTransforms allows taking the output of existing sources (Event/Entity/Join) and enriching them with 1 or more model outputs. This can be used in GroupBys, Joins, or hit directly via the fetcher. The GroupBy path allows for async materialization of model outputs to the online KV store for low latency serving. The fetcher path allows for on-demand model inference during online serving (at the cost of higher latency / more model inference calls).
Parameters:
sources— List of existing sources (Event/Entity/Join sources) to be enriched with model outputsmodels— List of Model objects that will be used for inference on the source datapassthrough_fields— Fields from the source that we want to passthrough alongside the model outputs - key_fields: List of tuples of (field_name, DataType) defining the schema of the key fields. If provided, creates a STRUCT schema that will be set as the ModelTransforms' keySchema. Example: [('user_id', DataType.STRING), ('session_id', DataType.STRING)]output_namespace— Namespace for the model outputtable_properties— Additional table properties for the model outputtags— Additional metadata tagsenvironments— List of environments where this ModelTransforms should be deployed/available. Defaults to ['prod']. Valid values: 'prod', 'canary' (case-insensitive).workflow_concurrency— Default maximum number of workflow steps Hub may allocate concurrently when a workflow is started from this ModelTransforms config. Request-level overrides take precedence.
OnlineStrategy
from ai.chronon.types import OnlineStrategy Operation
from ai.chronon.types import Operation| Name | Description |
|---|---|
MIN | Minimum value in the column |
MAX | Maximum value in the column |
FIRST | First non-null value of input column by time column |
LAST | Last non-null value of input column by time column |
APPROX_UNIQUE_COUNT | Approximate count of unique values using CPC (Compressed Probability Counting) sketch |
UNIQUE_COUNT | Exact count of unique values of the input column. Will store the set of items and can be expensive if the cardinality of the column is high. |
COUNT | Total count of non-null values of the input column |
SUM | Sum of values in the input column |
AVERAGE | Arithmetic mean of values in the input column |
VARIANCE | Statistical variance of values in the input column |
SKEW | Skewness (third standardized moment) of the distribution of values in input column |
KURTOSIS | Kurtosis (fourth standardized moment) of the distribution of values in input column |
HISTOGRAM | Full frequency distribution of values |
Operation
from ai.chronon.types import Operation| Name | Description |
|---|---|
COUNT | |
EXISTS | |
SUM | |
AVG | |
AVERAGE | |
MIN | |
MAX | |
LATEST | |
EARLIEST | |
COUNT_DISTINCT |
Query
from ai.chronon.types import Query
Query(
selects: Dict[str, str] = None,
wheres: List[str] = None,
start_partition: str = None,
end_partition: str = None,
time_column: str = None,
setups: List[str] = None,
mutation_time_column: str = None,
reversal_column: str = None,
partition_column: str = None,
partition_format: str = None,
partition_interval: Union[Window, str] = None,
partition_offset: Union[Window, str] = None,
partition_lag: Union[Window, str] = None,
sub_partitions_to_wait_for: List[str] = None,
time_partitioned: bool = None
) -> ttypes.QueryCreate a query object that is used to scan data from various data sources. This contains partition ranges, row level transformations and filtering logic. Additionally we also require a time_column for TEMPORAL events, mutation_time_column & reversal for TEMPORAL entities.
Parameters:
selects(List[str], optional) — Spark sql expressions with only arithmetic, function application & inline lambdas. You can also apply udfs see setups param below.:Example: { "alias": "built_in_function(col1) * my_udf(col2)", "alias1": "aggregate(array_col, 0, (acc, x) -> acc + x)" }See: https://spark.apache.org/docs/latest/api/sql/#built-in-functions When none, we will assume that no transformations are needed and will pick columns necessary for aggregations.wheres(List[str], optional) — Used for filtering. Same as above, but each expression must return boolean. Expressions are joined using AND.start_partition(str, optional) — From which partition of the source is the data valid from - inclusive. When absent we will consider all available data is usable.end_partition(str, optional) — Till what partition of the source is the data valid till - inclusive. Not specified unless you know for a fact that a particular source has expired after a partition and you should instead use another source after this partition.time_column(str, optional) — a single expression to produce time as ** milliseconds since epoch**.setups(List[str], optional) — you can register UDFs using setups ["ADD JAR YOUR_JAR", "create temporary function YOU_UDF_NAME as YOUR_CLASS"]mutation_time_column(str, optional) — For entities, with real time accuracy, you need to specify an expression that represents mutation time. Time should be milliseconds since epoch. This is not necessary for event sources, defaults to "mutation_ts"reversal_column(str, optional) — (defaults to "is_before") For entities with realtime accuracy, we divide updates into two additions & reversal. updates have two rows - one with is_before = True (the old value) & is_before = False (the new value) inserts only have is_before = false (just the new value). deletes only have is_before = true (just the old value). This is not necessary for event sources.partition_column(str, optional) — Specify this to override spark.chronon.partition.column set in teams.py for this particular query.sub_partitions_to_wait_for(List[str], optional) — Additional partitions to be used in sensing that the source data has landed. Should be a full partition string, such as `hr=23:00'partition_format(str, optional) — Date format string to expect the partition values to be in.partition_interval(Union[common.Window, str], optional) — Partition interval for the source table. Examples: "1d", "3h", "15m". Sub-daily partition values use "yyyy-MM-dd-HH-mm" unless partition_format is explicitly set.partition_offset(Union[common.Window, str], optional) — Offset from UTC midnight/epoch for the source partition grid. Examples: "1h" for partitions at 01:00, 04:00, ... with a "3h" partition_interval.partition_lag(Union[common.Window, str], optional) — Partition lag to apply when resolving dependencies. Examples: "1d", "3h", "15m".time_partitioned(bool, optional) — Indicates the source table uses a timestamp or date column for time-based filtering instead of traditional Hive-style string partitioning. When True, partition_column should be set to the timestamp/date column name. The engine will derive virtual partitions via MIN/MAX of that column and use timestamp-based WHERE clauses. Common for BigQuery, Snowflake, and Delta Lake tables that aren't Hive-partitioned.
Returns: A Query object that Chronon can use to scan just the necessary data efficiently.
ResourceConfig
from ai.chronon.types import ResourceConfig| Name | Description |
|---|---|
minReplicaCount | |
maxReplicaCount | |
machineType |
RolloutStrategy
from ai.chronon.types import RolloutStrategy| Name | Description |
|---|---|
rolloutType | |
validationTrafficPercentRamps | |
validationTrafficDurationMins | |
rolloutMetricThresholds |
ServingContainerConfig
from ai.chronon.types import ServingContainerConfig| Name | Description |
|---|---|
image | |
servingHealthRoute | |
servingPredictRoute | |
servingContainerEnvVars |
StagingQuery
from ai.chronon.types import StagingQuery
StagingQuery(
query: str,
version: Optional[int] = None,
output_namespace: Optional[str] = None,
table_properties: Optional[Dict[str, str]] = None,
setups: Optional[List[str]] = None,
engine_type: Optional[EngineType] = None,
dependencies: Optional[List[Union[TableDependency, Dict]]] = None,
tags: Optional[Dict[str, str]] = None,
offline_schedule: str = "@daily",
conf: Optional[ConfigProperties] = None,
env_vars: Optional[EnvironmentVariables] = None,
cluster_conf: ClusterConfigProperties = None,
step_days: Optional[int] = None,
recompute_days: Optional[int] = None,
additional_partitions: List[str] = None,
environments: Optional[List[str]] = None,
partition_interval: Optional[Union[Window, str]] = None,
partition_offset: Optional[Union[Window, str]] = None,
workflow_concurrency: Optional[int] = None
) -> ttypes.StagingQueryCreates a StagingQuery object for executing arbitrary SQL queries with templated date parameters.
Parameters:
query(str) — Arbitrary spark query that should be written with template parameters: -{{ start_date }}: Initial run uses start_date, future runs use latest partition + 1 day -{{ end_date }}: The end partition of the computing range -{{ latest_date }}: End partition independent of the computing range (for cumulative sources) -{{ max_date(table=namespace.my_table) }}: Max partition available for a given table These parameters can be modified with offset and bounds: -{{ start_date(offset=-10, lower_bound='2023-01-01', upper_bound='2024-01-01') }}setups(List[str]) — Spark SQL setup statements. Used typically to register UDFs.partition_column(str) —engine_type(int) — By default, spark is the compute engine. You can specify an override (eg. bigquery, etc.) Use the EngineType class constants: EngineType.SPARK, EngineType.BIGQUERY, etc.tags(Dict[str, str]) — Additional metadata that does not directly affect computation, but is useful for management.offline_schedule(str) — The offline schedule interval for batch jobs. Supports standard cron expressions, including regular sub-daily schedules. Examples: '@daily': Legacy format for midnight daily execution '0 2 * * *': Daily at 2:00 AM '0 */3 * * *': Every 3 hours '5/15 * * * *': Every 15 minutes with a 5 minute processing offset '@never': Explicitly disable offline schedulingpartition_interval(Optional[Union[common.Window, str]]) — Output partition interval for this StagingQuery. Examples: "1d", "3h", "15m". When set below daily and no partition format is supplied, Chronon uses "yyyy-MM-dd-HH-mm".partition_offset(Optional[Union[common.Window, str]]) — Offset from UTC midnight/epoch for the output partition grid. Defaults to zero (boundaries at midnight) and must be declared explicitly to move the grid; the cron fire phase is treated as a processing delay relative to the declared grid, never as a grid offset.conf(common.ConfigProperties) — Configuration properties for the StagingQuery.env_vars(common.EnvironmentVariables) — Environment variables for the StagingQuery.cluster_conf— Cluster configuration properties for the join.step_days(int) — The maximum number of days to process at oncedependencies(List[Union[TableDependency, Dict]]) — List of dependencies for the StagingQuery. Each dependency can be either a TableDependency object or a dictionary with 'name' and 'spec' keys.recompute_days(int) — Used by orchestrator to determine how many days are recomputed on each incremental scheduled run. Should be set when the source data is changed in-place (i.e. existing partitions overwritten with new data each day up to X days later) or when you want partially mature aggregations (i.e. a 7 day window, but start computing it from day 1, and refresh it for the next 6 days)environments(List[str]) — List of environments where this StagingQuery should be deployed/available. Defaults to ['prod']. Valid values: 'prod', 'canary' (case-insensitive).workflow_concurrency(int) — Default maximum number of workflow steps Hub may allocate concurrently when a workflow is started from this StagingQuery. Request-level overrides take precedence.
Returns: A StagingQuery object
TableDependency
from ai.chronon.types import TableDependencyDeclares an upstream table dependency for a StagingQuery. The orchestrator resolves the required upstream partition range from the query range and the offset / cutoff fields, then requires every partition in that range to be Filled before the step runs.
Resolved range: [max(query.start - start_offset, start_cutoff), min(query.end - end_offset, end_cutoff)]
:param offset:
DEPRECATED — prefer start_offset / end_offset. When set, applies
to both sides unless a specific start_offset / end_offset overrides
that side. Emits a DeprecationWarning.
:type offset: Optional[int]
:param start_offset:
Days subtracted from query.start to compute the upstream range start.
Combined with start_cutoff via max(query.start - start_offset, start_cutoff) — the cutoff acts as a floor. Defaults to 0 when unset,
except when only start_cutoff is set (no offset, no start_offset),
in which case the resolved startOffset is left None so the start
pins at start_cutoff regardless of the query range.
:type start_offset: Optional[int]
:param end_offset:
Days subtracted from query.end to compute the upstream range end.
Combined with end_cutoff via min(query.end - end_offset, end_cutoff) — the cutoff acts as a ceiling. Defaults to 0 when unset. Always
non-null in the resolved thrift: when end_cutoff is also unset, DependencyResolver.computeInputRange returns PartitionRange(_, null) and BatchNodeRunner.requiredEnd / the cumulative branch NPE downstream.
So end_cutoff can only clamp — it cannot pin the end the way start_cutoff can pin the start.
:type end_offset: Optional[int]
:param start_cutoff:
Fixed floor on the resolved range start (inclusive), as a partition
string (e.g. "2024-01-01"). Combined with start_offset as above.
:type start_cutoff: Optional[str]
:param end_cutoff:
Fixed ceiling on the resolved range end (inclusive), as a partition
string. Combined with end_offset as above.
:type end_cutoff: Optional[str]
| Name | Description |
|---|---|
resolved_spec | |
column | |
format | |
interval | |
offset | |
start | |
start | |
start | |
start | |
end | |
end | |
end | |
stacklevel | |
tableInfo | |
column | |
format | |
interval | |
offset | |
time_partitioned | |
table | |
startOffset | |
endOffset | |
startCutOff | |
endCutOff |
TimeUnit
from ai.chronon.types import TimeUnit| Name | Description |
|---|---|
MINUTES | |
HOURS | |
DAYS |
TrainingSpec
from ai.chronon.types import TrainingSpec| Name | Description |
|---|---|
trainingDataSource | |
trainingDataWindow | |
schedule | |
image | |
pythonModule | |
resourceConfig | |
jobConfigs |
Window
from ai.chronon.types import Window
Window(length: int, time_unit: TimeUnit) -> ttypes.Window