plugins.h_openlineage.OpenLineageAdapter
plugins.h_openlineage.OpenLineageAdapter
Install with pip install "apache-hamilton[openlineage]". The extra brings openlineage-python (the client) and openlineage-sql (the parser used to find the tables a query reads). openlineage-sql publishes no Windows wheel, so on Windows it is built from source during the install, which needs a Rust toolchain. Without the parser, SQL queries, and table names that are not plain identifiers, are reported without table datasets. The parser is used only with sql_dataset_identity="datasource"; the default identity needs neither it nor anything else beyond the client.
class hamilton.plugins.h_openlineage.OpenLineageAdapter(client: OpenLineageClient, namespace: str, job_name: str, sql_dataset_identity: Literal['legacy', 'datasource'] | None = None)[source]
This adapter emits OpenLineage events.
# create the openlineage client
from openlineage.client import OpenLineageClient
# write to file
from openlineage.client.transport.file import FileConfig, FileTransport
file_config = FileConfig(
log_file_path="/path/to/your/file",
append=False,
)
client = OpenLineageClient(transport=FileTransport(file_config))
# write to HTTP, e.g. marquez
client = OpenLineageClient(url="http://localhost:5000")
# create the adapter
adapter = OpenLineageAdapter(client, "my_namespace", "my_job_name")
# add to Hamilton
# import your pipeline code
dr = driver.Builder().with_modules(YOUR_MODULES).with_adapters(adapter).build()
# execute as normal -- and openlineage events will be emitted
dr.execute(...)Note for data lineage to be emitted, you must use the “materializer” abstraction to provide metadata. See https://hamilton.apache.org/concepts/materialization/. This can be done via the @datasaver() and @dataloader() decorators, or using the @load_from or @save_to decorators, as well as passing in data savers and data loaders via .with_materializers() on the Driver Builder, or via .materialize() on the driver object.
__init__(client: OpenLineageClient, namespace: str, job_name: str, sql_dataset_identity: Literal['legacy', 'datasource'] | None = None)[source]
Constructor. You pass in the OLClient.
Parameters:
- self
- client
- namespace
- job_name
- sql_dataset_identity – how SQL loader/saver datasets are identified.
"legacy"(the current default) names themnamespace+ bare table name, as earlier releases did."datasource"names them after the database they live in, per the OpenLineage naming conventions (seesql_datasets()), so lineage connects across jobs. Leaving it unset uses"legacy"and warns once, since the default will change to"datasource"in a future major release.
Returns:
post_graph_execute(run_id: str, graph: FunctionGraph, success: bool, error: Exception | None, results: dict[str, Any] | None)[source]
Emits a Run COMPLETE or FAIL event.
Parameters:
- run_id
- graph
- success
- error
- results
Returns:
post_node_execute(run_id: str, node_: Node, kwargs: dict[str, Any], success: bool, error: Exception | None, result: Any | None, task_id: str | None = None)[source]
Run Event: will emit a RUNNING event with updates on input/outputs.
A Job Event will be emitted for graph execution, and additional SQLJob facet if data was loaded from a SQL source.
A Dataset Event will be emitted if a dataloader or datasaver was used:
- input data set if loader
- output data set if saver
- appropriate facets will be added to the dataset where it makes sense.
TODO: attach statistics facets
Parameters:
- run_id
- node
- kwargs
- success
- error
- result
- task_id
Returns:
pre_graph_execute(run_id: str, graph: FunctionGraph, final_vars: list[str], inputs: dict[str, Any], overrides: dict[str, Any])[source]
Emits a Run START event. Emits a Job Event with the sourceCode Facet for the entire DAG as the job.
Parameters:
- run_id
- graph
- final_vars
- inputs
- overrides
Returns:
pre_node_execute(run_id: str, node_: Node, kwargs: dict[str, Any], task_id: str | None = None)[source]
No event emitted.
SQL datasets
SQL loaders and savers (@load_from.sql, @save_to.sql and the pandas SQL materializers) record the datasource they used (see SQL metadata and lineage). How the adapter names their datasets is set by sql_dataset_identity:
"legacy", the default: datasets are named as in earlier Hamilton releases, under the adapter’s job namespace with the baretable_name. As before, a query read containingSELECTproduces a dataset with no name and the query in the job’ssqlfacet, and other strings are used as the dataset name. Leaving the option unset warns once per adapter (aFutureWarning) because the default will change; pass"legacy"explicitly to keep these names without the warning."datasource": datasets are named after the datasource, following the OpenLineage naming conventions, and every physical table a query reads is reported. A report written by one job and read by another then resolves to the same dataset, and two tables with the same name in different databases stay distinct.
adapter = OpenLineageAdapter(client, "my_namespace", "my_job", sql_dataset_identity="datasource")The rest of this section describes the "datasource" identity:
| Dialect | Namespace | Name |
|---|---|---|
| PostgreSQL | postgres://{host}:{port} (port defaults to 5432) | {database}.{schema}.{table}; unquoted identifiers are folded to lower case, as the server does |
| SQLite | sqlite://{absolute file path}, with forward slashes on every platform (sqlite://C:/data/sales.db on Windows) | {table}. A table in an attached database (reporting.orders in the SQL, or the writer’s schema) is named in the attached file’s namespace, which is known only for a standard-library sqlite3 connection; otherwise it is left out. |
Aliases and common table expressions are not reported as tables. Each dataset carries a dataSource facet with the namespace; the schema facet from dataframe_metadata is attached only when the node maps to a single table. The job keeps the sql facet with the statement’s text; a node that read or wrote a table by name has no statement to report. The exception is a table read by a name that the metadata files as a query (one containing SELECT, or whose first word is select or with, such as SELECT_LOG or select-log): it can’t be told apart from a statement, so the name is reported as the job’s SQL and no input dataset is emitted for it. Give such tables plain names, or read them with a query. A read by another name that is not a plain identifier (daily revenue) is emitted as that table when no SQL statement starts at it; one starting with table is left out, since Postgres’ TABLE t is a statement openlineage-sql does not know.
Whatever cannot be fully identified is left out and logged as a warning from the hamilton.plugins.h_openlineage logger, never guessed. The following cases are left out:
- an unknown datasource (in-memory SQLite, an unsupported connection object, legacy two-argument metadata)
- an unsupported dialect
- a table whose schema cannot be determined
- a statement
openlineage-sqlcannot parse - a missing
openlineage-sqlinstall
A failure inside the conversion is logged with its exception type and the run event is still emitted without datasets. The node itself has already succeeded and is never failed by lineage.
Migrating to the datasource identity
The default stays "legacy" for now and will become "datasource" in a future major release. Switching changes every SQL dataset’s namespace and name (for example from my_namespace + daily_revenue to postgres://warehouse.example:5432 + analytics.reporting.daily_revenue). A lineage backend shows the new names as new datasets: history recorded under the old names does not connect to them, and nothing is rewritten automatically. Some datasets also stop appearing:
- loaders and savers whose metadata has no
source(custom functions using the two-argument helper, in-memory databases) - any of the unidentifiable cases above, including every SQL query and every table name read or written that is not a plain identifier wherever
openlineage-sqlis not installed
To migrate:
- Install
openlineage-sqlwhere you run Hamilton.apache-hamilton[openlineage]includes it; on Windows it is built from source, which needs a Rust toolchain. Pin it in locked environments. - Run the pipeline once with
sql_dataset_identity="datasource"against a test backend, or read the events withFileTransport, and note the new namespace and name of each dataset. Check thehamilton.plugins.h_openlineagewarnings for anything left out. - In your lineage backend, move what is keyed to the old names (ownership, tags, alerts, policies, saved queries) to the new ones, or link old and new datasets where the backend supports it.
- Pass
sql_dataset_identity="datasource"in production.
To stay on the current names, pass sql_dataset_identity="legacy". It silences the warning.
Reuse the conversion
sql_datasets() is the boundary another integration (an orchestrator provider, a custom adapter) can call on metadata Hamilton produced, with no OpenLineageClient, Driver or database connection involved:
from hamilton.plugins.h_openlineage import sql_datasets
lineage = sql_datasets(node_result_metadata) # the dict a SQL loader/saver returned
lineage.inputs # list[openlineage.client.event_v2.Dataset]
lineage.outputs
lineage.notes # why anything was left outhamilton.plugins.h_openlineage.sql_datasets(sql_metadata: dict[str, Any], operation: Literal['read', 'write'] | None = None) → SqlDatasets[source]
Converts Hamilton SQL metadata into OpenLineage datasets. Emits nothing, opens nothing.
This is the reusable boundary for other integrations (e.g. an orchestrator provider): feed it the sql_metadata produced by hamilton.io.utils.get_sql_metadata() (either the whole metadata dict or its sql_metadata entry) and get datasets named per the OpenLineage naming conventions:
- PostgreSQL: namespace
postgres://{host}:{port}, name{database}.{schema}.{table}. Unquoted identifiers are folded to lower case, as the server does. - SQLite: namespace
sqlite://{absolute file path}(forward slashes on every platform), name{table}. A table in an attached database ({schema}.{table}in the SQL, or the writer’sschema) is named in the attached file’s namespace; it is left out when that file cannot be known.
Queries are parsed with openlineage-sql; every physical table read appears in inputs and every table written in outputs (aliases and common table expressions are not tables). A bare table name is placed by operation (read/write), taken from the metadata unless given here. What a write names is its table when it is a plain identifier; otherwise it is parsed, and is the table written unless it parses as a statement that names tables (then the tables it writes are used, if any). Without a parse it is left out. A read string that is not a plain name is parsed too, and is the table read only when no statement starts at it (daily revenue); otherwise it is SQL, and without a parse it is left out. Schema precedence: explicit in the SQL, then the writer’s schema, then the connection’s default schema. Anything that cannot be fully identified is left out and explained in notes rather than guessed — an unknown datasource, an unsupported dialect, a missing schema, a parse error or a missing openlineage-sql install.
class hamilton.plugins.h_openlineage.SqlDatasets(inputs: list[Dataset], outputs: list[Dataset], notes: list[str], query: str | None = None)[source]
Datasets resolved from Hamilton SQL metadata. notes explains anything left out.
query is the SQL statement the metadata recorded, or None when it named a table (a written name may be filed under the metadata’s query) or cannot be told apart from one. A table read by a name the metadata files as a query (one containing SELECT, or whose first word after comments is select or with in any case) cannot be told apart from a statement, so its name is reported here and no input dataset is emitted for it.
评论
登录后参与评论
KnowForge