DatabricksCopyIntoOperator¶
Use the DatabricksCopyIntoOperator to import
data into Databricks table using COPY INTO
command.
Using the Operator¶
Operator loads data from a specified location into a table using a configured endpoint. The only required parameters are:
table_name- string with the table namefile_location- string with the URI of data to loadfile_format- string specifying the file format of data to load. Supported formats areCSV,JSON,AVRO,ORC,PARQUET,TEXT,BINARYFILE.One of
sql_endpoint_name(name of Databricks SQL endpoint to use) orhttp_path(HTTP path for Databricks SQL endpoint or Databricks cluster).
Other parameters are optional and could be found in the class documentation.
Examples¶
Importing CSV data¶
An example usage of the DatabricksCopyIntoOperator to import CSV data into a table is as follows:
# Example of importing data using COPY_INTO SQL command
import_csv = DatabricksCopyIntoOperator(
task_id="import_csv",
databricks_conn_id=connection_id,
sql_endpoint_name=sql_endpoint_name,
table_name="my_table",
file_format="CSV",
file_location="abfss://container@account.dfs.core.windows.net/my-data/csv",
format_options={"header": "true"},
force_copy=True,
)
DatabricksCopyIntoAssetOperator¶
Use DatabricksCopyIntoAssetOperator
when the COPY INTO target is a Unity Catalog table that downstream Dags schedule on.
It accepts every DatabricksCopyIntoOperator argument plus a required unity_table.
DatabricksCopyIntoOperator itself declares no assets.
unity_table is a UnityTableIdentity
with host, catalog, schema, and table. All four are required and static.
Jinja in any field raises ValueError. host follows the Databricks connection rule, so
https://my-workspace.cloud.databricks.com/ becomes my-workspace.cloud.databricks.com.
The hostname, catalog, schema, and table are normalized to lowercase because their names
are case-insensitive. Different capitalization therefore produces the same asset URI.
unity_table.to_asset() returns the databricks://host/catalog/schema/table asset.
from airflow.providers.databricks.assets.databricks import UnityTableIdentity
from airflow.providers.databricks.operators.databricks_sql import DatabricksCopyIntoAssetOperator
users = UnityTableIdentity(
host="my-workspace.cloud.databricks.com",
catalog="main",
schema="default",
table="users",
)
load_users = DatabricksCopyIntoAssetOperator(
task_id="load_users",
sql_endpoint_name="my-endpoint",
file_location="/Volumes/main/default/landing/users.csv",
file_format="CSV",
table_name="main.default.users",
unity_table=users,
)
Outlets¶
When you omit outlets, the operator sets outlets=[unity_table.to_asset()] at parse time.
When you pass outlets, including outlets=[], the operator keeps your value.
To add extra outlets, pass the Unity table asset among them, as in
outlets=[users.to_asset(), other_asset].
Templated table names¶
table_name stays templated. unity_table is not templated and fixes the asset at parse time.
Before running SQL, the operator resolves the rendered table_name and compares it with unity_table.
The table and connection hostname comparisons are case-insensitive; the SQL keeps the supplied casing.
A three-part name
catalog.schema.tableis used as is.A two-part name
schema.tabletakes the catalog from thecatalogargument.A one-part name
tabletakes the catalog and schema from thecatalogandschemaarguments.
The operator never guesses the workspace default catalog or schema. If a part is missing or the
resolved table differs from unity_table, the task raises ValueError and no SQL runs.
Dynamic task mapping¶
A mapped task does not run the operator constructor when the Dag is parsed, so it gets no
automatic outlet. Pass outlets=[users.to_asset()] in partial(). All mapped instances share
one unity_table, so each rendered table_name must still resolve to that table.
Sources and other operators¶
file_location may be a Unity Catalog volume path such as /Volumes/main/default/landing.
The volume is the COPY INTO source. It is not an outlet. The outlet is the target table.
DatabricksSQLStatementsOperator does not
infer assets from SQL. Pass outlets=[users.to_asset()] for the tables your statements write.
Databricks job and pipeline operators, such as
DatabricksRunNowOperator, do not infer
table assets either. Pass outlets explicitly.