Lecture Overview & Objectives
Learning Objectives (as given)
- Describe the purpose and benefits of Lakeflow Connect for scalable data ingestion into Databricks.
- Identify the different types of connectors, including Standard and Managed connectors.
- Explain various data ingestion techniques such as batch, incremental batch, and streaming.
- Select the appropriate ingestion method based on data and use case requirements.
- Review the key benefits of Delta tables and the Medallion Architecture for data management and analytics.
- Explain the need for managed connectors when ingesting data from enterprise databases and SaaS applications beyond cloud object storage.
- Describe the benefits of Lakeflow Connect Managed Connectors including simplified setup, UI-driven configuration, and fully managed infrastructure.
- Compare the SaaS and database ingestion architectures used by Lakeflow Connect Managed Connectors.
- Describe Partner Connect as an alternative for data sources without a native managed connector.
- Describe additional Databricks features that extend data integration, sharing, and collaboration capabilities.
- Identify key Databricks Marketplace components and explore shared data assets.
- Explain the purpose of
MERGE INTOfor applying updates, inserts, and deletes to existing Delta tables. - Describe
MERGE INTOclauses including matched updates, matched deletes, and not matched inserts. - Write a
MERGE INTOstatement to merge source data into a target Delta table.
1. Purpose & Benefits of Lakeflow Connect
Objective 1
Lakeflow Connect is Databricks’ unified data ingestion platform — “simple and efficient connectors to ingest data from local files, popular enterprise applications, databases, cloud storage, message buses, and more” into the Databricks Data Intelligence Platform.
Lakeflow Connect is one of three pillars of the broader Lakeflow product family:
- Connect — efficient ingestion connectors (this lecture)
- Spark Declarative Pipelines (SDP) — accelerated ETL development (the current name for what was previously Delta Live Tables / DLT — see gotcha below)
- Jobs — reliable orchestration for analytics and AI workloads
All three run on Databricks’ “industry leading data processing engine” (Apache Spark + Structured Streaming), sit under Unity Catalog for unified governance, and write to Delta Lake, Parquet, or Iceberg for optimized storage.
- Sources
- On-prem systems
- Databases
- Cloud storage
- SaaS applications
- Event streams
- Lakeflow Connect
- Ingestion connectors
- Data Intelligence Platform
- Unified governance
- Efficient compute & storage
- Use cases
- Streaming analytics
- Data science & ML
- Data sharing
- BI & reporting
- Breadth of sources — on-prem systems, databases, cloud object storage, SaaS applications, and event streams all ingest through one platform.
- Unity Catalog governance — ingested data lands directly into Unity Catalog–governed tables, so access control and lineage are consistent from the moment data arrives.
- Powers downstream consumers — the same ingested data feeds streaming analytics, data science/ML, data sharing, and BI & reporting without separate pipelines per consumer.
- Efficient compute & storage — managed connectors run on serverless compute and use incremental reads/writes to avoid reprocessing unchanged data.
What are the three components of the Lakeflow product family, and which one does this lecture cover?
Connect (ingestion), Spark Declarative Pipelines / SDP (ETL), and Jobs (orchestration). This lecture covers Connect.
Does Lakeflow Connect ingest data into Unity Catalog–governed tables by default?
Yes — Lakeflow Connect is governed by Unity Catalog, so ingested data lands in governed tables with consistent access control and lineage.
2. Connector Types: Upload, Standard & Managed
Objective 2
Databricks gives you three broad ways to get data into the platform:
- Upload Files — upload local files directly to Databricks, upload a file to a volume, or create a table straight from a local file. Good for one-off, small, manual loads.
- Standard Connectors — connect to cloud object storage, Kafka, and other sources, supporting all three ingestion methods (batch, incremental batch, streaming).
- Managed Connectors — fully-automated connectors for SaaS applications and databases, with incremental reads/writes built in.
| Upload files | Standard connectors | Managed connectors | |
|---|---|---|---|
| Sources | Local files | Cloud object storage, Kafka and other sources | SaaS applications and databases |
| Ingestion methods | Manual, one-off | Batch, incremental batch and streaming | Incremental reads and writes, automated |
| Setup | Upload through the UI, to a volume or straight into a table | You configure it: more work, more control | Minimal: UI, API or SDK |
| Best for | Small, manual loads | Custom file or stream ingestion | Faster, scalable, more cost-efficient ingestion from enterprise systems |
Customizable connectors for cloud object storage and message buses (e.g. Kafka). Require more configuration than managed connectors but offer greater flexibility. Support all three ingestion methods: batch, incremental batch, and streaming.
Fully-automated connectors built on Lakeflow pipelines, with source-specific authentication, CDC (change data capture), edge-case handling, long-term API maintenance, automated retries, and automated schema evolution baked in. Require minimal configuration (UI, APIs, or SDKs).
Per Databricks docs, managed connectors fall into six categories:
- Database connectors (CDC) — e.g. MySQL, PostgreSQL, SQL Server
- Query-based connectors — direct database queries without CDC configuration
- SaaS connectors — e.g. Salesforce, HubSpot, Jira, Workday
- File source connectors — e.g. Google Drive, SharePoint
- Streaming connectors — e.g. RabbitMQ, event streams
- Community connectors — open-source, community-built
On the first run a managed connector ingests all selected data; subsequent runs ingest only what changed — “faster, scalable, and more cost-efficient” than re-pulling everything each time.
You need to ingest data from Salesforce with minimal configuration and automatic schema evolution. Standard or Managed connector?
Managed connector — Salesforce is a SaaS source, and managed connectors handle authentication, CDC, and schema evolution automatically.
You need to ingest from a Kafka topic with full control over batch vs. streaming configuration. Standard or Managed?
Standard connector — Kafka is a message bus source supported by standard connectors, which require more configuration but give you the flexibility to choose the ingestion method.
Sources
3. Ingestion Techniques: Batch, Incremental Batch, Streaming
Objective 3
| Method | What each run ingests | Common techniques |
|---|---|---|
| Batch | All data, every time | SQL CREATE TABLE AS SELECT (CTAS); Python spark.read.load() |
| Incremental batch | Only new data; already-loaded records are skipped | COPY INTO; Auto Loader (spark.readStream with a timed trigger); CREATE OR REFRESH STREAMING TABLE |
| Streaming | Data continuously as it’s generated, in short micro-batches | Auto Loader with a continuous trigger; Declarative Pipelines in continuous mode |
- Loads data as batches of rows into Databricks, often on a schedule.
- Traditional batch ingestion re-processes all records every time it runs — no memory of what was already loaded.
- Common techniques: SQL
CREATE TABLE AS SELECT(CTAS); Pythonspark.read.load(). - Best for smaller, one-time or ad hoc datasets where reprocessing everything isn’t costly.
- Ingests (appends) only new data — previously loaded records are skipped automatically.
- Faster and more resource-efficient than full batch since it processes less data per run.
- Common techniques: SQL
COPY INTO; Pythonspark.readStream(Auto Loader with a timed trigger); Declarative PipelinesCREATE OR REFRESH STREAMING TABLE.
- Continuously loads data rows or small batches of rows as they are generated, so they can be queried in near real-time.
- Micro-batch processes small batches at very short, frequent intervals — this is how Structured Streaming implements “continuous” processing under the hood.
- Common techniques: Python
spark.readStream(Auto Loader with a continuous trigger); Declarative Pipelines in continuous trigger mode.
Primary ways to ingest data — comparison
| Feature | CTAS + spark.read (Batch) | COPY INTO (Incremental Batch) | Auto Loader (Incremental: Batch or Streaming) |
|---|---|---|---|
| Use cases | Best for smaller datasets | Ideal for thousands of files | Scales to millions+ files/hour; backfills billions of files |
| Syntax/Interface | Python (spark.read) or SQL (CTAS) |
SQL | Python (spark.readStream), SQL Declarative Pipelines (CREATE OR REFRESH STREAMING TABLE), streaming tables in Databricks SQL |
| Idempotency | No | Yes | Yes |
| Schema evolution | Manual or inferred on read | Supported with options | Automatically detects & evolves schema; handles new columns as they appear |
| Latency | High | Moderate (scheduled) | Low or high, depending on configuration |
| Ease of use | Simple | Simple, SQL-based | Intermediate to advanced |
| Summary | Best for one-time, ad hoc ingestion; can be scheduled to always reprocess all data | Simple, repeatable incremental file ingestion; great for scheduled jobs/pipelines | Best for near real-time or incremental ingestion, with high automation and scalability |
What's the key difference between batch and incremental batch ingestion?
Batch re-processes all records every run. Incremental batch ingests only new data, automatically skipping previously loaded records.
Which ingestion technique(s) are idempotent?
COPY INTO and Auto Loader are idempotent. CTAS / spark.read is not.
You need to ingest millions of files per hour with automatic schema evolution. Which tool?
Auto Loader — built to scale to millions+ files/hour with automatic schema detection and evolution.
Sources
4. Choosing the Right Ingestion Method
Objective 4
Use the comparison table in §3 to reason about trade-offs along these axes:
- Data volume — small/ad hoc → CTAS; thousands of files → COPY INTO; millions+ → Auto Loader.
- Ingestion frequency / latency needs — one-time or infrequent → batch; scheduled incremental → incremental batch; near real-time → streaming (Auto Loader with continuous trigger, or Declarative Pipelines).
- Idempotency requirements — if re-running a job must not duplicate data, avoid CTAS/spark.read alone; prefer COPY INTO or Auto Loader.
- Schema volatility — frequently changing source schemas favor Auto Loader’s automatic schema evolution over manual CTAS handling.
- Source type & governance needs — SaaS apps/databases with governance and minimal ops overhead → Managed Connectors; cloud storage/Kafka needing custom control → Standard Connectors.
A vendor sends flat files to a cloud storage bucket at an unpredictable, high volume — sometimes millions of files. Which method?
Auto Loader — it’s purpose-built to scale to millions+ files per hour with automatic schema handling, and can run as either incremental batch or streaming.
You have a small, static reference CSV to load once into a table. Which method?
CTAS / spark.read — simplest option, appropriate for smaller one-time datasets where idempotency and schema evolution aren’t concerns.
5. Delta Lake Review
Objective 5
Delta Lake is the open-source storage layer underlying Databricks’ lakehouse. Delta tables hold all three data quality levels of the Medallion Architecture — raw (bronze), cleaned (silver), and aggregated (gold) — as one consistent format, so you don’t need separate storage systems per stage.
- ACID transactions — safe concurrent reads/writes, no partial or corrupted data.
- Time travel — query or restore previous versions of a table for auditing and rollback.
- Unified batch + streaming — the same Delta table can be a sink for both batch and Structured Streaming writes/reads.
- Schema enforcement & evolution — rejects mismatched writes by default, but can be configured to evolve as new columns appear.
Name two Delta Lake features that make it suitable for both rollback and regulatory audit trails.
ACID transactions (guarantee consistent, non-corrupted writes) and time travel (query/restore prior table versions).
6. Medallion Architecture
Objective 5
- Ingest
- Batch
- Streaming
- Bronze
- Raw data
- Silver
- Cleaned data
- Gold
- Business-level aggregates
- Consumers
- BI & reporting
- ML & AI
- Streaming analytics
- Ingest Data — data enters Delta Lake via batch, streaming, or both; this is the starting point for all processing.
- Process & Improve Data Quality — data is incrementally refined as it moves through layers, each stage improving structure, quality, and usability.
- Bronze Layer (Raw Data) — stores raw, unprocessed data from multiple sources; the foundation for all downstream processing.
- Silver Layer (Cleaned Data) — data is cleaned, transformed, and enriched; produces structured, analysis-ready datasets.
- Gold Layer (Business-Ready Data) — curated, aggregated data optimized for reporting, BI, and advanced analytics.
| Layer | Contains | Typical consumer |
|---|---|---|
| Bronze | Raw, unvalidated data in original format(s); single source of truth for auditing | Data engineers |
| Silver | Cleaned, deduplicated, validated, schema-enforced; structured & analysis-ready | Analysts / downstream pipelines |
| Gold | Highly aggregated, business-level, dimensionally modeled | BI dashboards, executives, ML/AI, streaming analytics |
Which medallion layer is the "single source of truth" for auditing, and why?
Bronze — it retains the raw, unmodified data exactly as ingested, before any cleaning or transformation, so it can always be replayed or audited against the original source.
Why does progressing data through bronze → silver → gold help both BI and ML/AI teams share one dataset?
Because each layer is a Delta table with consistent ACID guarantees and Unity Catalog governance — BI and ML/AI both read from the same governed Gold tables rather than maintaining separate, potentially inconsistent copies.
7. Lab Environment & Unity Catalog Basics
Section 1 topic
Unity Catalog metadata for every database object is organized in a 3-level namespace: catalog.schema.object (e.g. table, view, volume, function, model).
- Catalog — top-level container, often mirroring an org unit or environment (e.g. prod vs dev).
- Schema — sits inside a catalog; holds tables, views, volumes, models, and functions.
- Object — the actual data asset (table, view, volume, etc.).
You can set a default catalog and schema for a session so you don’t have to fully qualify every object reference (e.g. write my_table instead of my_catalog.my_schema.my_table).
A schema is a container for several distinct object types, each governed the same way by Unity Catalog but holding a different kind of asset:
Catalog
└── Schema
├── Tables → structured data
├── Views → saved query, no storage of its own
├── Volumes → arbitrary files
├── Functions → registered UDFs
└── Models → registered ML model artifacts| Object type | What it holds |
|---|---|
| Tables | Structured data organized into rows and columns (e.g. Delta tables). |
| Views | A saved SQL query over one or more tables/views — read-only, and stores no data of its own; it re-runs the query on each access. |
| Volumes | Governed access to arbitrary, non-tabular files (e.g. images, PDFs, model checkpoints) sitting in cloud storage. |
| Functions | Reusable, executable logic — registered user-defined functions (UDFs) and stored procedures. |
| Models | Versioned (or unversioned) registered ML model artifacts and their metadata. |
| Command | Purpose |
|---|---|
USE CATALOG / USE SCHEMA |
Sets the current default catalog/schema for the session, so unqualified object references resolve against it. Note: switching catalog resets the current schema back to default. |
IDENTIFIER(...) |
SQL-injection-safe way to parameterize identifiers (table/column/schema/function names) in a statement, e.g. CREATE TABLE IDENTIFIER(mytab)(c1 INT) instead of unsafe string concatenation. |
SHOW (e.g. SHOW TABLES, SHOW SCHEMAS, SHOW VOLUMES) |
Lists objects of a given type within the current or specified catalog/schema. |
DESCRIBE / DESCRIBE TABLE |
Shows metadata (columns, types, comments) for a given object. |
DESCRIBE SCHEMA EXTENDED |
Shows extended metadata about a schema, including location and properties, beyond the basic DESCRIBE SCHEMA output. |
DESCRIBE VOLUME |
Shows metadata for a Unity Catalog volume. |
LIST '/Volumes/catalog/schema/volume/...' |
Lists the files stored inside a Unity Catalog volume path. |
TABLE |
Shorthand SQL statement to query/select all rows from a table (equivalent to SELECT * FROM). |
What are the three levels of the Unity Catalog namespace, in order?
Catalog → Schema → Object (table/view/volume/function/model).
What happens to the current schema when you run USE CATALOG to switch catalogs?
It resets to default — you need to run USE SCHEMA again after switching catalogs if you want a non-default schema active.
Why would you use IDENTIFIER(...) instead of string-concatenating a table name into SQL?
To safely parameterize the identifier and avoid SQL-injection risk, especially when the name comes from a variable or external input.
Sources
8. Code Syntax Reference (SQL)
Databricks SQL snippets from the lecture’s notebooks, organized by what they do.
Session context — catalog & schema
SELECT current_catalog(), current_schema(); -- show the active catalog & schema
USE CATALOG <catalog_name>; -- switch the session's default catalog
USE SCHEMA <schema_name>; -- switch the session's default schema
current_catalog()/current_schema()— return the catalog/schema currently active for the session, as strings.USE CATALOG/USE SCHEMA— set the session’s default catalog/schema so unqualified object references resolve against it (see §7). Switching catalog resets the current schema todefault.
Exploring objects
SHOW TABLES; -- list tables in the current catalog/schema
SHOW VOLUMES; -- list volumes in the current catalog/schema
SHOW <object_type> lists objects of that type within the current (or a specified) catalog/schema — same family as SHOW SCHEMAS, SHOW CATALOGS, SHOW FUNCTIONS, etc. Full command table is in §7.
Reading schema/column metadata — information_schema
SELECT 'new_hires_ctas' AS table_name, COUNT(*) AS column_count
FROM information_schema.columns
WHERE table_schema = 'get_started_de' AND table_name = 'new_hires_ctas'
UNION ALL
SELECT 'employees', COUNT(*)
FROM information_schema.columns
WHERE table_schema = 'get_started_de' AND table_name = 'employees';
information_schema.columns— a SQL-standard system view that describes every column of every table/view you have access to: which schema and table it belongs to, its column name, data type, and more. Instead of eyeballingDESCRIBE TABLEone table at a time, you can query it like any other table.- This query uses the same
UNION ALLcomparison trick as §8’sVERSION AS OFexample — running two filteredCOUNT(*)s side by side, one per table, to compare column counts betweennew_hires_ctasandemployeesat a glance (handy for sanity-checking that two tables you expect to line up actually do). WHERE table_schema = '...' AND table_name = '...'— filters the metadata down to one specific table, sinceinformation_schema.columnsotherwise lists columns for every table/view you can see.
If you query information_schema.columns with no catalog prefix, which catalog’s metadata do you get?
Your session’s current catalog (§7) — information_schema is scoped per catalog, so switching catalogs with USE CATALOG changes what it shows unless you explicitly qualify it (e.g. my_catalog.information_schema.columns).
Two colleagues run the exact same information_schema query and get different results. Why might that be, even against the same catalog?
information_schema is privilege-aware — it only returns objects the querying user actually has Unity Catalog permission to see, so different access levels can produce different results.
Exploring files & reading raw files
LIST 'volumes/../'; -- list files/objects at a volume path
SELECT * FROM read_files('<filepath>'); -- read files using Databricks' default format detection
LIST '<path>'— lists files stored at a Unity Catalog volume path.read_files(path [, option => value, ...])— a table-valued function that reads files directly into a queryable table shape. It auto-detects the file format (JSON, CSV, XML, TEXT, BINARYFILE, PARQUET, AVRO, ORC) and infers a unified schema across files when no format is specified.
Previewing raw file content before ingestion
SELECT * FROM text.`<file_path>`
LIMIT 5;
- Spark SQL lets you reference files directly as a table by prefixing the path with a format name:
<format>.`<path>`— backticks around the path, not quotes.text.`<path>`reads the file(s) at that path using the text data source: one row per line, in a singleSTRINGcolumn namedvalue, with zero parsing. - This is a quick way to eyeball raw file content before committing to a parsing strategy — check for a header row, confirm the shape looks like JSON/CSV, spot base64-encoded fields — the same reconnaissance step
read_files(above) skips straight past. LIMIT 5caps it to the first 5 lines, same as any other query.
Creating tables — batch ingestion example
CREATE TABLE IF NOT EXISTS <table_name>
AS SELECT * FROM read_files('<filepath>');
-- Creates the table only if it doesn't already exist, populated from the ingested file(s).
-- This is the CTAS batch-ingestion pattern from §3.
This is CREATE TABLE AS SELECT (CTAS) — the exact batch-ingestion technique named in §3’s comparison table. It creates table_name with a schema inferred from the query, and populates it in one step from read_files(...).
What’s wrong with CREATE TABLES IF NOT EXISTS my_table SELECT * FROM read_files('/path')?
Two errors: TABLES should be singular TABLE, and it’s missing the AS keyword before SELECT. Corrected: CREATE TABLE IF NOT EXISTS my_table AS SELECT * FROM read_files('/path').
Which function reads files directly into a table shape with automatic format detection?
read_files(path) — supports JSON, CSV, XML, TEXT, BINARYFILE, PARQUET, AVRO, and ORC, auto-detecting format and inferring a unified schema when no format option is given.
CTAS with parsing options — controlling how read_files() reads a file
CREATE TABLE IF NOT EXISTS current_employees_ctas AS
SELECT ID, FirstName, Country, Role
FROM read_files(
'/Volumes/dbacademy/get_started_de/myfiles/employees.csv',
format => 'csv',
header => true,
inferSchema => true
);
- Same CTAS pattern as above, but instead of
SELECT *, it projects only specific columns (ID, FirstName, Country, Role) —read_filescan be queried like any other table source, so normal column selection, filtering, joins, etc. all work on top of it. format => 'csv'— tellsread_filesexplicitly to parse the source as CSV, instead of relying on auto-detection.header => true— treats the first row of the file as column headers rather than data.inferSchema => true— asks Spark to infer each column’s data type from the file contents, instead of reading every column in as a string.- These are passed as
key => valueoptions directly inside theread_files(...)call — the general pattern isread_files(path, option_key => option_value, ...), and you can stack as many as you need. - The explicit column list here also has a second purpose:
read_filesautomatically adds a_rescued_datacolumn to its output (see gotcha below) — selecting the four named columns instead of*is what excludes it from the final table.
- The whole point is to rescue, not drop: instead of silently discarding data that doesn’t fit the schema (or erroring out the whole ingest), the offending value(s) get preserved in this one column so nothing is lost.
- It’s stored as a JSON-formatted string — per the official docs, “a JSON blob with the rescued columns and the source file path of the record.” So it’s a STRING column holding a JSON object, not a native struct — you’d use a JSON function (e.g.
get_json_object, orfrom_jsonwith a schema) to pull specific values back out of it if needed. - Specifically, the docs name three triggers for a field landing in
_rescued_data: it’s absent from the provided/inferred schema, its value’s data type doesn’t match the schema, or its column name has a case mismatch against the schema’s field names. - A row with no mismatches has an empty/no rescued content for that row — matching the general “only populated when something didn’t fit” design of the feature (see the syntax note below on confirming the exact
NULLwording).
In read_files(path, format => 'csv', header => true, inferSchema => true), what does inferSchema => true do?
Tells Spark to detect each column’s actual data type from the file’s contents, rather than reading every column in as a string.
Can you SELECT specific columns (not *) directly from a read_files(...) call inside a CTAS?
Yes — read_files produces a normal queryable table shape, so column projection, filters, and joins all work on top of it just like any other table/view.
Why might SELECT * FROM read_files(...) return a column you didn’t expect, and how do you avoid it?
read_files automatically adds a _rescued_data column to hold any data that didn’t match the schema. Select named columns instead of * to exclude it, or set schemaEvolutionMode => 'none' to turn it off entirely.
What format is the data stored in inside _rescued_data, and what three things can trigger a field landing there?
A JSON-formatted string (a JSON blob with the rescued columns and the source file path). It’s triggered by a field being absent from the schema, having a data type that doesn’t match the schema, or having a column-name case mismatch against the schema’s field names.
Does ingesting a perfectly clean file (no schema mismatches at all) mean the _rescued_data column disappears from the table?
No — the column is still part of the table’s schema; it’s just NULL (empty) for every row where nothing needed rescuing. It only gets removed if you explicitly turn the feature off (e.g. schemaEvolutionMode => 'none') or select named columns that exclude it.
File metadata column — _metadata
-- Select everything plus the hidden metadata struct:
SELECT *, _metadata FROM read_files('<dir_path>');
-- Or just the fields you need, e.g. for lineage/audit columns:
SELECT
*,
_metadata.file_name AS source_file_name,
_metadata.file_path AS source_file_path,
_metadata.file_modification_time AS source_modified_at
FROM read_files('<dir_path>');
_metadata— a hiddenSTRUCTcolumn available on any file-based data source (CSV, JSON, Parquet, etc. — whether read viaread_files, plainspark.read/DataFrameReader, Auto Loader, orCOPY INTO). Like_rescued_data, it isn’t shown bySELECT *unless you ask for it by name — it has to be explicitly selected to appear.- Fields inside the struct:
file_path(STRING),file_name(STRING, with extension),file_size(LONG, bytes),file_modification_time(TIMESTAMP),file_block_start(LONG) andfile_block_length(LONG) — the last two describe which byte-range of the file was read. - This is exactly the “additional context, lineage information, and data governance” use case: instead of just the data, you can carry along where each row came from — which source file, when it was last modified, how big it was — as ordinary queryable columns on the target table, e.g. for auditing which upstream file introduced a given row.
- Dot-notation (
_metadata.file_name) pulls out one field from the struct, and — as in the second query above — you can alias it to a friendlier column name that actually lands in your table.
What does the _metadata column give you that plain file columns don’t?
Per-row provenance about the source file itself — its path, name, size, and last-modified time — useful for lineage, auditing, and data governance, rather than the file’s actual data content.
Why won’t COPY INTO target FROM '<path>' ... let you select _metadata directly?
COPY INTO’s simple FROM '<path>' form doesn’t expose it — you need the subquery form instead: COPY INTO target FROM (SELECT *, _metadata FROM '<path>') FILEFORMAT = ....
What happens if a source file already has a real column called _metadata?
That column name is reserved — queries return the data source’s own _metadata column instead of the file metadata struct, silently shadowing the feature rather than raising an error.
Auto Loader via SQL — CREATE OR REFRESH STREAMING TABLE
CREATE OR REFRESH STREAMING TABLE
catalog.schema.table
SCHEDULE EVERY 1 HOUR
AS
SELECT * FROM STREAM read_files(
'<dir_path>',
format => '<file_type>'
)
CREATE OR REFRESH STREAMING TABLE catalog.schema.table AS SELECT ...— the third way to define a table’s contents, alongside CTAS and the explicit-schema form above. Declares a streaming table: a Delta table designed to incrementally process a growing input, handling each source row only once. This is the SQL/Declarative Pipelines side of Auto Loader named in §3/§4’s comparison table.STREAM read_files('<dir_path>', format => '<file_type>')— wraps the sameread_filestable-valued function from earlier in §8 in aSTREAMmodifier, turning a normal batch read into an incremental streaming source that only picks up new files dropped at the path.SCHEDULE EVERY 1 HOUR— tells the pipeline to refresh (check for and ingest new files) once every hour, rather than running continuously. Valid units areHOUR/HOURS(1–72),DAY/DAYS(1–31), andWEEK/WEEKS(1–8) — there’s also aCRONform for more precise schedules. OmittingSCHEDULEentirely runs the table continuously instead of on a fixed interval.CREATE OR REFRESH— combines “create if it doesn’t exist” and “update the definition if it does” into one statement, so re-running this same DDL after editing it (e.g. changing the schedule or source path) updates the streaming table in place.
What does wrapping read_files(...) in STREAM change about how it behaves?
It turns a normal one-time batch read into an incremental streaming source — instead of reading everything at the path every time, it only picks up files that weren’t already processed on a previous run.
What does SCHEDULE EVERY 1 HOUR control, and what happens if you omit SCHEDULE entirely?
It sets how often the streaming table refreshes (checks for and ingests new data) — once per hour, in this case. Omitting SCHEDULE runs the table continuously instead of on a fixed interval.
Creating an empty table — explicit column definitions
CREATE TABLE IF NOT EXISTS practice_copyinto (
ID INT,
FirstName STRING,
Country STRING,
Role STRING
);
- A second flavor of
CREATE TABLE, alongside the CTAS form above: instead ofAS SELECT ... FROM read_files(...), you spell out atable_specification—(col_name col_type, col_name col_type, ...)— directly. This creates an empty Delta table with the given schema and zero rows, ready to be loaded separately afterward. IF NOT EXISTS— same semantics as the CTAS version: no error if the table already exists. It also won’t alter the schema of an existing table if one is already there with a different shape.- No
USINGclause is given, so the table defaults to Delta format — per the official docs, “IfUSINGis omitted, the default isDELTA.” INT,STRING— standard Databricks SQL column types (there’s alsoBIGINT,DOUBLE,BOOLEAN,DATE,TIMESTAMP, and more).
What’s the key structural difference between CREATE TABLE ... AS SELECT ... and CREATE TABLE table_name (col1 type1, col2 type2)?
CTAS infers the schema from a query and populates the table with the query’s results in one step. The explicit-schema form defines columns yourself and creates an empty table with no rows — you load data into it separately afterward.
What table format does practice_copyinto use, given the statement has no USING clause?
DELTA — that’s the default table format whenever USING is omitted from CREATE TABLE.
Modifying data — INSERT, UPDATE, DELETE
INSERT INTO <table_name> VALUES (...); -- add new row(s) by literal values
INSERT INTO <table_name> SELECT * FROM <other_table>; -- add new row(s) from a query
UPDATE <table_name>
SET <column> = <value>
WHERE <condition>; -- modify existing row(s) that match
DELETE FROM <table_name>
WHERE <condition>; -- remove row(s) that match
INSERT INTO— appends new rows, either as literalVALUESor from the result of aSELECTquery. Supports Delta schema enforcement/evolution, and aBY NAMEoption to match columns by name instead of position.UPDATE ... SET ... WHERE— modifies existing rows in place. TheWHEREclause is optional — omit it and every row gets updated.DELETE FROM ... WHERE— removes rows that match the condition. Also optional — omitWHEREand it deletes every row in the table.
Upserting with MERGE INTO
| Clause | When it fires | Allowed actions |
|---|---|---|
WHEN MATCHED |
The key exists in both source and target | UPDATE or DELETE |
WHEN NOT MATCHED [BY TARGET] |
The row is only in the source (it’s new) | INSERT |
WHEN NOT MATCHED BY SOURCE |
The row is only in the target (it’s gone upstream) | UPDATE or DELETE |
A single, atomic operation that applies updates, inserts, and deletes from a source into a target Delta table in one pass — supporting schema enforcement or schema evolution depending on how it’s configured. This makes it a natural fit for slowly changing dimensions (SCD), incremental loads, and change data capture (CDC) scenarios, where a batch of upstream changes (new rows, updated rows, deleted rows) needs to be applied to a target table together.
MERGE INTO target_table AS target
USING source_table AS source
ON target.id = source.id
WHEN MATCHED AND source.is_deleted = true THEN
DELETE
WHEN MATCHED THEN
UPDATE SET *
WHEN NOT MATCHED THEN
INSERT *;
MERGE INTO target USING source ON <merge_condition>— compares asourcetable/view/query against atargetDelta table row-by-row using the merge condition (usually a key match liketarget.id = source.id), then applies one or moreWHENclauses depending on whether each row matched.WHEN MATCHED— the merge condition found a row in both target and source. Action is eitherUPDATE SET ...(update the matched target row) orDELETE(remove it).UPDATE SET *updates every column by name from the source;UPDATE SET col = expr, ...updates specific columns.WHEN NOT MATCHED [BY TARGET]— the source row has no match in the target (i.e. it’s new). Action isINSERT— eitherINSERT *(insert all source columns by name) orINSERT (col1, ...) VALUES (...)for specific columns.WHEN NOT MATCHED BY SOURCE— the reverse: a target row has no match in the source (i.e. it’s missing from the incoming data). Action isUPDATE SET ...orDELETE— useful for marking/removing target rows that no longer exist upstream.- Multiple
WHENclauses can be chained (as in the example above — a conditional delete, then a general update, then an insert), each optionally narrowed further withAND <condition>. Clauses are evaluated in order, and only the first matching clause per row is applied.
Worked example — applying delete, update & insert in one MERGE INTO
Before: target_table
| users | status | |
|---|---|---|
| peter | [email protected] |
current |
| zebi | [email protected] |
current |
Incoming changes: source_table
| users | status | |
|---|---|---|
| peter | [email protected] |
delete |
| zebi | [email protected] |
update |
| samarth | [email protected] |
new |
After the merge: target_table
| users | status | |
|---|---|---|
| zebi | [email protected] |
update |
| samarth | [email protected] |
new |
MERGE INTO target_table target
USING source_table source -- declaring source & target
ON target.id = source.id -- merge condition
WHEN MATCHED AND source.status = 'update' THEN -- 1st MATCHED clause
UPDATE SET
target.email = source.email,
target.status = source.status
WHEN MATCHED AND source.status = 'delete' THEN -- 2nd MATCHED clause
DELETE
WHEN NOT MATCHED THEN -- NOT MATCHED clause (else)
INSERT (id, first_name, email, sign_up_date, status)
VALUES (source.id, source.first_name, source.email, source.sign_up_date, source.status);
target_tablestarts with two users (peter,zebi), bothstatus = 'current'.source_tablecarries three rows describing what should happen to each:peter→'delete',zebi→'update'(new email),samarth→'new'(doesn’t exist in target yet).- The merge condition
ON target.id = source.idmatchespeterandzebi(both exist in both tables) but notsamarth(only in source) — that’s what makessamarthfall through to theWHEN NOT MATCHEDclause instead of aWHEN MATCHEDone. - Two separate
WHEN MATCHEDclauses, each narrowed withAND source.status = ..., route matched rows to different actions based on the value in the source row — not just whether a match exists.zebi(status'update') hits the first clause and gets itsemail/statusupdated;peter(status'delete') hits the second clause and is removed entirely. - Result:
target_tableends up withzebi(updated email, statusupdate) andsamarth(newly inserted, statusnew) —peteris gone. All three outcomes (update, delete, insert) applied from one source table in a single atomic statement.
What’s the difference between WHEN NOT MATCHED and WHEN NOT MATCHED BY SOURCE in a MERGE INTO statement?
WHEN NOT MATCHED (equivalently WHEN NOT MATCHED BY TARGET) fires for source rows with no matching target row — the action is an INSERT. WHEN NOT MATCHED BY SOURCE fires for target rows with no matching source row — the action is an UPDATE or DELETE.
You need to insert new rows, update changed rows, and delete rows flagged as removed — all from a single CDC source table, in one statement. What's the right tool?
MERGE INTO — it’s the single statement that can combine conditional insert, update, and delete actions (via WHEN MATCHED / WHEN NOT MATCHED clauses) against a Delta target table in one pass.
In the worked example, why does samarth get inserted instead of hitting one of the WHEN MATCHED clauses, even though source_table has a status column for it just like the other two rows?
The WHEN MATCHED vs. WHEN NOT MATCHED branching is driven entirely by the merge condition (target.id = source.id), not by the status column’s value. Since no row in target_table has samarth’s id, it falls through every WHEN MATCHED clause and lands in WHEN NOT MATCHED, regardless of what its status column says.
Two WHEN MATCHED clauses both apply to the same matched row (say, one with AND source.status = 'update' and one with no extra condition at all). Which one runs?
Only the first one that matches — WHEN clauses are evaluated in the order they’re written, and the first clause whose condition is satisfied is the one applied per row; later clauses are skipped for that row.
Table history — DESCRIBE HISTORY
DESCRIBE HISTORY <table_name>; -- show the audit log of every write to this Delta table
DESCRIBE HISTORY table_name returns provenance information for every write to a Delta table — version number, timestamp, the user who ran it, and the operation (e.g. WRITE, UPDATE, DELETE, MERGE). This is the audit trail underpinning Delta Lake’s time travel feature (§5) — you can query or restore any prior version shown here.
Querying a previous version — VERSION AS OF
-- Query an older version directly:
SELECT * FROM new_employees VERSION AS OF 0;
SELECT * FROM new_employees TIMESTAMP AS OF '2026-01-01T00:00:00.000Z';
-- Compare current vs. a previous version, e.g. row counts before/after a change:
SELECT 'Current' AS version, COUNT(*) AS row_count FROM new_employees
UNION ALL
SELECT 'Version 0', COUNT(*) FROM new_employees VERSION AS OF 0;
VERSION AS OF <n>— queries the table exactly as it looked at versionn(the version numbers come straight out ofDESCRIBE HISTORY).TIMESTAMP AS OF '<timestamp>'— same idea, but pinned to a point in time instead of a version number.- Because it’s just a modifier on the table reference, you can use it anywhere a table can appear in a query — including a
UNION ALLlike the example above, which is a handy way to eyeball what changed (e.g. row counts, or specific rows) between “now” and an earlier version, without altering the table itself.
What happens if you run DELETE FROM my_table; with no WHERE clause?
Every row in the table is deleted — WHERE is optional, and omitting it applies the operation to all rows.
You need to insert a row where it doesn't already exist, or update it where it does, in one statement. Which command?
MERGE INTO (an upsert) — not plain INSERT INTO or UPDATE alone.
What does DESCRIBE HISTORY let you see, and what Delta Lake feature does it power?
It shows the version, timestamp, user, and operation for every write to the table — the audit trail that powers time travel (querying/restoring prior table versions).
Does SELECT * FROM my_table VERSION AS OF 3 change what’s currently stored in my_table?
No — it only reads version 3 for that query. To actually revert the table’s stored state, you’d need RESTORE TABLE my_table TO VERSION AS OF 3.
Sources
- read_files table-valued function — Databricks on AWS
- current_catalog function — Databricks on AWS
- current_schema function — Databricks on AWS
- CREATE TABLE [USING] — Databricks on AWS
- CREATE STREAMING TABLE — Databricks on AWS
- INSERT — Databricks on AWS
- UPDATE — Databricks on AWS
- DELETE FROM — Databricks on AWS
- MERGE INTO — Databricks on AWS
- Upsert into a Delta Lake table using merge — Databricks on AWS
- DESCRIBE HISTORY — Databricks on AWS
- Work with table history — Databricks on AWS
- RESTORE — Databricks on AWS
- Batch read options (DataFrameReader) — Databricks on AWS
- Information schema — Databricks on AWS
- File metadata column — Databricks on AWS
- Configure schema inference and evolution in Auto Loader — Databricks on AWS
9. Code Syntax Reference (PySpark)
Python/PySpark notebook snippets, organized the same way as §8’s SQL reference.
Querying a table safely — spark.sql(), display(), and error handling
# Run this after completing the Upload UI steps above.
try:
display(spark.sql("SELECT * FROM current_employees_ui"))
except Exception as e:
if "TABLE_OR_VIEW_NOT_FOUND" in str(e):
print("The table 'current_employees_ui' doesn't exist yet.")
print("Complete the Upload UI steps above, then re-run this cell.")
else:
raise e
spark.sql("...")— runs a SQL string from Python and returns the result as a DataFrame, the PySpark equivalent of a%sqlcell (§8). Any SQL from §8 can be run this way inside a Python notebook.display(...)— the Databricks notebook function for rendering a DataFrame as a scrollable, interactive table in the cell output, with a built-in option to turn it into a chart. It’s a Databricks notebook feature, not a plain PySpark/Python function — different from a DataFrame’s own.show()(plain text output) or Python’sprint().try / except Exception as e— wraps the query so a missing table doesn’t crash the whole cell/notebook run with an ugly stack trace.if "TABLE_OR_VIEW_NOT_FOUND" in str(e):— checks the exception’s text for this specific Databricks/Spark error condition name, so theexceptblock reacts only to “table doesn’t exist yet” and not to some other, unrelated failure.else: raise e— re-raises anything that isn’t the expected missing-table error, so a real bug (bad SQL, permissions issue, etc.) still surfaces loudly instead of being silently swallowed.
What’s the difference between display(df) and df.show() in a Databricks notebook?
display() is a Databricks notebook feature that renders an interactive, scrollable table with a one-click option to add a chart; .show() is plain PySpark and prints a static text table to the console with no interactivity.
Why does this code check "TABLE_OR_VIEW_NOT_FOUND" in str(e) instead of just catching and swallowing every exception?
So it only handles the specific, expected case (the table doesn’t exist yet) and re-raises (raise e) anything else — an unrelated bug shouldn’t be silently hidden behind a generic “table doesn’t exist” message.
Is TABLE_OR_VIEW_NOT_FOUND specific to this notebook, or a real Databricks/Spark error condition?
It’s an official error condition name — raised whenever Databricks can’t resolve a referenced table or view (missing object, typo, wrong catalog/schema, or missing permissions).
Batch ingestion with DataFrameReader/DataFrameWriter — the Python equivalent of CTAS
#1. Read the parquet file from the volume into a Spark DataFrame
df = (spark
.read
.format("parquet")
.load("/Volumes/<catalog path>/<volume path>/<file name>")
)
#2. Write DataFrame to a Delta table (overwrite if exists)
(df
.write
.mode("overwrite")
.saveAsTable(f"<catalog_path>.{DA.schema_name}.<file_name>")
)
#3. Read and view table
bronze_table = spark.table(f"<catalog_path>.{DA.schema_name}.<file_name>")
bronze_table.display()
spark.read.format("parquet").load(path)— the plain PySpark DataFrameReader, reading a file straight into a DataFrame. This is Batch ingestion (§3/§4), same category as CTAS — just the Python API instead of SQL, and without going throughread_files..write.mode("overwrite").saveAsTable(name)— the DataFrameWriter counterpart: writes the DataFrame out as a managed Delta table (Delta is the default format, as in §8). Together, steps 1–2 are the Python/DataFrameReader equivalent of a single SQLCREATE TABLE ... AS SELECT ...statement — read, then materialize as a table — just split into two explicit calls instead of one.mode("overwrite")— replaces the table’s contents (and, per the docs, its schema too — an overwrite isn’t required to match the existing table’s columns) if it already exists, rather than erroring or appending. This is more permissive than CTAS’sIF NOT EXISTS, which refuses to touch an existing table at all.f"<catalog_path>.{DA.schema_name}.<file_name>"— an f-string building the full three-level Unity Catalog name (§7) from a literal catalog, a variable schema name (DA.schema_name, from the Databricks Academy lab helper), and a literal table name.spark.table(name)— the DataFrameReader shortcut for reading an existing table by name into a DataFrame (equivalent tospark.sql("SELECT * FROM name"), just without writing SQL).
What are steps 1–2 of this snippet the Python/DataFrameReader equivalent of, in SQL?
A single CREATE TABLE ... AS SELECT ... (CTAS) statement — reading a file and materializing it as a table, just split into an explicit read (spark.read...load()) and write (.write...saveAsTable()) instead of one SQL statement.
Does spark.read.format("parquet").load(path) add a _rescued_data column automatically, the way read_files does?
No — plain DataFrameReader reads don’t add it automatically. You’d need to opt in explicitly with .option("rescuedDataColumn", "_rescued_data") to get the same safety net.
COPY INTO via spark.sql() — incremental batch loading
result = spark.sql("""
COPY INTO practice_copyinto
FROM '/Volumes/dbacademy/get_started_de/myfiles'
FILEFORMAT = CSV
FORMAT_OPTIONS ('header' = 'true', 'inferSchema' = 'true')
""")
result.display()
COPY INTO target_table FROM <location> FILEFORMAT = <format> [FORMAT_OPTIONS(...)]— a SQL command that loads new files from a source location into an existing Delta table. This is the Incremental Batch ingestion technique named in §3/§4’s comparison table, and it’s the table-loading step that pairs with the empty table created just above.spark.sql("""...""") — same pattern as this section’s first entry: runs a SQL string from Python. Triple quotes let the multi-lineCOPY INTOstatement be written across several lines without escaping newlines.result = spark.sql(...)—COPY INTOreturns a DataFrame summarizing the operation (e.g. counts of rows/files inserted vs. skipped) when its result is captured like this, rather than left unassigned.result.display()— renders that summary DataFrame as an interactive table, samedisplay()covered above.FORMAT_OPTIONS ('header' = 'true', 'inferSchema' = 'true')— same idea asread_files’sheader/inferSchemaoptions in §8: treat row 1 of each CSV as column headers, and infer each column’s data type from file contents instead of reading everything as a string.
What makes COPY INTO safe to re-run on a schedule against the same source folder, unlike CTAS?
COPY INTO tracks which files it has already loaded into the target table and skips them on subsequent runs, even if they were modified since — it’s idempotent/incremental. CTAS has no such tracking; running the CREATE statement again does nothing once the table exists, and there’s no repeatable “load only what’s new” behavior.
Why is the COPY INTO statement wrapped in triple quotes inside spark.sql("""...""")?
Triple-quoted Python strings can contain literal newlines, so the multi-line SQL statement can be written across several lines for readability without needing to escape them.
What does capturing the result (result = spark.sql(...)) of a COPY INTO call give you?
A DataFrame summarizing the load operation — e.g. how many rows/files were inserted vs. skipped — the same kind of object any spark.sql() call returns, so it can be passed to display() like any other DataFrame.
COPY INTO bronze_table
FROM '<file_path>'
FILEFORMAT = parquet
COPY_OPTIONS ('mergeSchema' = 'true');
A worked example of the mergeSchema case above: if files landing at <file_path> start including columns that bronze_table doesn’t have yet, COPY_OPTIONS ('mergeSchema' = 'true') lets the table’s schema evolve to pick up those new columns automatically, instead of the load failing or silently dropping the extra data. Without it, a schema mismatch between source files and target table causes COPY INTO to error.
Auto Loader — streaming ingestion with readStream/writeStream
(spark
.readStream
.format("cloudFiles")
.option("cloudFiles.format", "json")
.option("cloudFiles.schemaLocation", "<checkpoint_path>")
.load("/Volumes/catalog/schema/files")
.writeStream
.option("checkpointLocation", "<checkpoint_path>")
.trigger(processingTime="5 seconds")
.toTable("catalog.database.table")
)
spark.readStream.format("cloudFiles")— the Python entry point for Auto Loader, the third ingestion technique from §3/§4."cloudFiles"is the special Structured Streaming source name that activates Auto Loader’s incremental file-detection engine, instead of a normal one-time batch read..option("cloudFiles.format", "json")— tells Auto Loader what format the source files are in (also supports CSV, Parquet, Avro, text, binaryFile, ORC — same format family asread_filesin §8)..option("cloudFiles.schemaLocation", ...)— a path where Auto Loader persists the schema it infers, so it remembers the schema (and any evolution) across restarts instead of re-inferring from scratch every run..load("/Volumes/catalog/schema/files")— the source directory Auto Loader watches for new files..writeStream.option("checkpointLocation", ...)— a separate path (conceptually distinct fromschemaLocation, though both are often placed under the same parent directory) where Structured Streaming tracks processing progress/state, enabling exactly-once, resume-safe streaming..trigger(processingTime="5 seconds")— runs as a micro-batch stream, checking for new files every 5 seconds — this is the “Streaming” row from §3’s comparison (near real-time via frequent small batches), as opposed to a one-off scheduled run..toTable("catalog.database.table")— the three-level Unity Catalog target table (§7) this stream continuously writes into.
What does setting .format("cloudFiles") on spark.readStream actually activate?
Auto Loader — Structured Streaming’s incremental file-detection source, which tracks which files have already been ingested instead of re-reading everything on each run.
What’s the difference between cloudFiles.schemaLocation and checkpointLocation?
schemaLocation persists the inferred/evolved schema across runs; checkpointLocation persists streaming progress state (what’s already been processed), enabling exactly-once, resume-safe processing. They’re different paths serving different purposes.
How do you make an Auto Loader stream run once over all currently available files and then stop, instead of streaming continuously?
Use .trigger(availableNow=True) instead of a fixed processingTime trigger — this makes Auto Loader run in incremental-batch mode rather than continuous streaming, e.g. for a scheduled job.
Sources
- PySpark basics — Databricks on AWS
- SparkSession.sql — Databricks on AWS
- TABLE_OR_VIEW_NOT_FOUND error condition — Databricks
- DataFrameReader.format — Databricks on AWS
- DataFrameWriter.saveAsTable — Databricks on AWS
- COPY INTO — Databricks on AWS
- Common data loading patterns using COPY INTO — Databricks on AWS
- Read and write CSV files — Databricks on AWS
- What is Auto Loader? — Databricks on AWS
- Common data loading patterns with Auto Loader — Databricks on AWS
- Ingest data from cloud object storage — Databricks on AWS
10. Ingestion & Governance: Explicit Schemas, Rescued Data & File Metadata
The conceptual home for ingestion’s data-quality/governance features. Full syntax mechanics for _rescued_data live in §8’s CTAS section, and for _metadata in §8’s dedicated subsection — this section covers why explicit schemas matter, and newer worked examples combining schema + rescued data + troubleshooting.
Why define an explicit schema?
- Reduces inferred-schema inconsistency risk — especially for semi-structured sources like JSON or CSV, where Spark has to guess types from sampled data. An explicit schema removes the guessing entirely: every column’s name and type is exactly what you declared, run after run.
- Faster parsing/loading — Spark can apply the declared types and structure immediately, rather than first scanning data to work out what the schema probably is.
- Lower compute overhead at scale — schema inference isn’t free: Auto Loader’s own docs describe sampling “the first 50 GB or 1000 files that it discovers, whichever limit is crossed first” just to infer a schema before real ingestion even starts. On a large or frequently-refreshed source, skipping that sampling pass every time by supplying the schema yourself is a real, measurable savings.
- This is also why the
inferSchema/inferColumnTypesnaming ambiguity flagged elsewhere in these notes (§8) stops mattering once you provide an explicitschema— there’s no inference step left to worry about.
Defining an explicit schema with Auto Loader / read_files
SELECT *
FROM read_files(
'<file_path>',
format => 'csv',
sep => '|',
header => true,
schema => 'order_id INT, email STRING, transactions_timestamp BIGINT',
rescuedDataColumn => '_rescued_data' -- names the rescued-data column
);
schema => 'order_id INT, email STRING, transactions_timestamp BIGINT'— a DDL-format schema string, passed as a single option, tellingread_filesexactly what columns to expect and what type each one is. No inference happens once this is supplied.sep => '|'— the field delimiter option for CSV (default is,); here the source file is pipe-delimited instead of comma-delimited.rescuedDataColumn => '_rescued_data'— names the column that holds anything that doesn’t fit the schema above. Any row with an extra field, a missing field, a type mismatch, or a column-name case mismatch against this schema gets that offending data rescued into_rescued_datainstead of the row being dropped or the whole read failing — same mechanic already covered in §8, just now paired with an explicit schema instead of an inferred one.
Diagnosing and fixing a misread header from _rescued_data
SELECT
CAST(_rescued_data:_c0 AS BIGINT) AS order_id,
*
FROM read_files(
'<file_path>',
format => 'csv',
sep => '|',
header => true,
schema => 'order_id INT, email STRING, transactions_timestamp BIGINT',
rescuedDataColumn => '_rescued_data'
);
- The scenario: after an initial ingest (or a check of the raw
_rescued_datavalues), you notice entries keyed likec0/_c0inside_rescued_datainstead of the real column name. That positional-style key (_c0,_c1, …) is Spark’s standard fallback naming for a column it couldn’t map to a name — e.g. the header row didn’t parse the way the schema expected, so the first column’s value got rescued under a generic positional key rather than being matched toorder_id. _rescued_data:_c0— the colon (:) operator, Databricks SQL’s JSON path-extraction syntax. It works directly on a STRING column containing valid JSON (exactly what_rescued_datais), pulling out the value at key_c0without needing a separatefrom_jsonparse step.CAST(... AS BIGINT)— the colon operator returns the extracted value typed the same as the JSON it came from (effectively untyped/string-like), so it’s cast to the correct target type — here recovering what should have been theorder_idcolumn.AS order_id, *— the recovered, correctly-typed value is aliased back to its intended column name and selected alongside everything else (*, which still includes the original columns and the still-present_rescued_datastruct) — effectively a manual repair of the one column that failed to parse.
Name the three benefits of defining an explicit schema instead of letting Spark infer one.
Reduced risk of inferred-schema inconsistencies (especially for semi-structured sources like JSON/CSV), faster parsing/loading since Spark applies the declared types immediately, and lower compute overhead at scale since there’s no sampling pass needed to infer the schema.
Does setting rescuedDataColumn => '_rescued_data' turn on the rescued-data feature?
No — a rescued-data column is provided by default. This option only controls what that column is named, which matters mainly to avoid a collision with a real column already called _rescued_data.
What does a key like _c0 inside _rescued_data tell you?
That Spark couldn’t map that value to a named column from the schema — it fell back to a positional placeholder name instead, usually because the header row wasn’t parsed the way the schema/format options expected.
What does the : operator do in _rescued_data:_c0, and why is CAST(...) needed around it?
The colon operator extracts the value at JSON path _c0 directly from the _rescued_data string (which holds valid JSON). The extracted value isn’t typed as the target column should be, so CAST(... AS BIGINT) converts it to the correct type before treating it as a real order_id.
Sources
11. Semi-Structured Data: JSON, STRUCT & VARIANT
Lakeflow Connect ingests semi-structured sources like JSON too (§2’s exam objective on semi-structured/unstructured ingestion) — this section covers JSON’s basic shape and the three ways to work with a JSON-formatted string column once it’s landed in a table.
JSON basics — objects, keys & values
{
"name": "John Doe",
"age": 35,
"address": { "city": "Anytown", "state": "CA" },
"children": [
{ "name": "Owen", "age": 10 },
{ "name": "Eva", "age": 8 }
]
}
- A JSON object (the outer
{ ... }) is a set of key/value pairs — same idea as a Python dict or a row of named columns. - A value’s type can be a string (
"name": "John Doe"), number ("age": 35), boolean, array ("children": [...]), object ("address": {...}), or null. - Objects can be flat (every value a simple string/number/boolean) or nested (a value is itself an object or an array of objects) — the
addressfield above is a nested object, andchildrenis an array of nested objects. - How deep and irregular this nesting gets is exactly what determines how hard the source is to work with downstream — a flat JSON file behaves a lot like a CSV with named columns; deeply nested/variable JSON is where STRUCT and VARIANT (below) start to matter.
Three ways to work with a JSON-formatted string column
- STRING — store the JSON as raw text, untouched. Simplest option, but you can’t query individual fields without a JSON function on every access, and there’s no schema enforcement at all.
- STRUCT — parse the JSON into a fixed, defined schema (via
from_json, below), giving you a proper typed, nested column you can dot into directly (e.g.col.address.city). Requires knowing/declaring the schema up front, and a schema change on the source breaks it. - VARIANT — a semi-structured type built to store any JSON shape flexibly, without a fixed schema, while still being efficiently queryable (path syntax like
col:address:city). Per the official docs, “the improved read and write performance for variant allows it to replace native Spark complex types such as structs and arrays in some use cases” — Databricks explicitly recommendsVARIANTover plain JSON strings for semi-structured data going forward.
Step 1 — deriving the schema with schema_of_json
SELECT schema_of_json('<sample-json-string>');
-- Returns the inferred schema as a STRING, e.g.:
-- STRUCT<name: STRING, age: BIGINT, address: STRUCT<city: STRING, state: STRING>, children: ARRAY<STRUCT<name: STRING, age: BIGINT>>>
- Instead of hand-writing a DDL schema string (the way §10’s
read_filesexample did),schema_of_json(jsonStr)takes one example JSON-formatted string and automatically derives the matching schema definition from it — useful when the shape is complex or you don’t want to type it out by hand. - The nested
addressobject becomes a nestedSTRUCT<...>, and thechildrenarray of objects becomesARRAY<STRUCT<...>>— the returned schema mirrors the JSON’s own nesting exactly, object-for-object and array-for-array. - The result is returned as a plain
STRING— a schema definition, not parsed data — which is exactly the shapefrom_json(next step) expects as itsschemaargument.
Step 2 — parsing JSON with from_json
SELECT from_json(json_col, '<json-struct-schema>') AS struct_column
FROM table;
from_json(jsonStr, schema [, options])— takes a STRING column holding JSON text plus a schema (a DDL string, or aschema_of_json(...)call directly), and returns a new column typed as aSTRUCTmatching that schema.- This is Step 1 + Step 2 as a pipeline: derive the schema once with
schema_of_jsonon a representative sample, then feed that schema intofrom_jsonto actually parse every row’s JSON string into a real typed/nested column — from here on, you can dot into fields directly (e.g.struct_column.address.city) instead of using JSON string functions. - Column names in the schema must match the JSON’s own key names exactly (case-sensitive) — a mismatch there behaves the same as the schema-mismatch cases already covered for
read_files/_rescued_datain §8/§10, just via a different function.
A more complex example — object holding an array of structs
Using the same JSON as above, the matching STRUCT schema is:
STRUCT<
name: STRING,
age: INT,
address: STRUCT<city: STRING, state: STRING>, -- one nested object
children: ARRAY<STRUCT<name: STRING, age: INT>> -- a list of objects
>
- The walkthrough breaks building a STRUCT schema for nested JSON into five steps: (1) define the schema for the JSON string overall, (2) wrap it in
STRUCT<...>, (3) declare simple fields’ types (name: STRING,age: INT), (4) declare a nested-object field as its ownSTRUCT<...>(address), and (5) declare an array-of-objects field asARRAY<STRUCT<...>>(children). - The key distinction highlighted here: a single nested object (
address, one city/state pair) is justSTRUCT<...>, but a list of similarly-shaped objects (children, multiple name/age pairs) needs theARRAY<>wrapper around theSTRUCT<>— getting this wrong (e.g. declaringchildrenas a plainSTRUCTinstead ofARRAY<STRUCT>) is a common schema-mismatch source, and would land those rows’childrendata in_rescued_data(§8/§10) instead of parsing correctly.
In JSON, what's the difference between a nested object and an array of objects, in terms of the values they can hold?
A nested object (like address) holds exactly one set of key/value pairs. An array of objects (like children) holds zero or more sets of key/value pairs, each with the same general shape, as a list.
How would you declare a schema field for a JSON key that holds a list of {name, age} objects?
ARRAY<STRUCT<name: STRING, age: INT>> — the ARRAY<> wrapper is required around the STRUCT<> because it’s a list of objects, not a single object.
What are the three ways to work with a JSON-formatted string column, and what's the key trade-off of each?
STRING (raw text, simplest, no schema enforcement, needs a JSON function per access), STRUCT (parsed into a fixed schema via from_json, typed and dot-accessible, but breaks if the source schema changes), and VARIANT (flexible schema-free storage with efficient path-based querying and better read/write performance than STRING/STRUCT in many cases, at the cost of not supporting clustering/partitioning/Z-order or direct comparison/grouping/ordering).
What does schema_of_json return, and what’s it typically used for immediately afterward?
It returns a STRING containing the inferred schema definition (as DDL-style STRUCT/ARRAY syntax) for a sample JSON string. That schema string is then typically passed straight into from_json(jsonCol, schema) to actually parse a JSON column into a typed STRUCT.
Base64-encoded values in semi-structured sources
Streaming/event sources (Kafka-style key/value columns, for example) commonly base64-encode their payloads. Base64 turns arbitrary bytes into a safe, printable ASCII string, which avoids corruption or encoding issues in transport (nulls, special characters, binary-safety) — but it also means the value isn’t human-readable until it’s decoded.
unbase64(expr) reverses this: it takes a base64-encoded STRING and returns the decoded value as BINARY — not a string. To get back to human-readable text, that binary result still needs to be CAST to STRING: CAST(unbase64(col) AS STRING).
-- bronze_table_raw -> bronze_table_decoded: decode base64 key/value columns into a new "bronze layer 2" table
CREATE OR REPLACE TABLE bronze_table_decoded AS
SELECT
CAST(unbase64(key) AS STRING) AS decoded_key,
<field_1>,
<field_2>,
CAST(unbase64(value) AS STRING) AS decoded_value
FROM bronze_table_raw;
This is the “bronze layer 2” pattern: bronze_table_raw holds the data exactly as ingested — base64-encoded key/value columns and all — and this step decodes those columns into a new, separate table (decoded_key / decoded_value) without touching or overwriting the raw source. decoded_value is now a JSON-formatted STRING column, ready for the STRING, STRUCT, and VARIANT flattening approaches below.
Flattening JSON via STRING conversion
Once a column holds a JSON-formatted STRING (like decoded_value above), the colon-path operator already introduced in §10 (col:key) can pull out its top-level — and nested — fields directly, with no conversion to STRUCT or VARIANT at all.
CREATE OR REPLACE TABLE bronze_string_flattened_1 AS
SELECT
<field_1>,
<field_2>,
decoded_key,
decoded_value:device,
decoded_value:traffic_source,
decoded_value:geo, -- can itself be a nested JSON-formatted string
decoded_value:items -- can be a nested array of JSON-formatted strings
FROM bronze_table_decoded;
- Each
decoded_value:<key>extracts one top-level field from the JSON string — same colon-path syntax as §10. - A colon-path result can itself hold JSON:
decoded_value:geois called out here as potentially another nested JSON-formatted string, anddecoded_value:itemsas a nested array of JSON-formatted strings. Going one level deeper (e.g.decoded_value:geo:city) works the same way. - This is the “go deeper as needed” version of the STRING approach — no upfront schema, just chain colon-paths further for whatever nested field you need, at the cost of no type enforcement and needing a colon-path (or JSON function) at every level.
Flattening JSON via STRUCT conversion
Converting the same decoded_value column to a STRUCT follows the two-step pattern already covered above: derive the schema with schema_of_json, then apply it with from_json.
SELECT schema_of_json('{...}') AS schema;
CREATE OR REPLACE TABLE bronze_table_struct AS
SELECT
* EXCEPT (decoded_value), -- all columns except the JSON-formatted string
from_json(decoded_value, schema) AS value
FROM bronze_table_decoded;
* EXCEPT (decoded_value)selects every column frombronze_table_decodedexceptdecoded_valueitself — shorthand for “keep everything else, I’m replacing this one column.”from_json(decoded_value, schema)parses that excluded column into a new, typedvalueSTRUCT column, using the schema string produced byschema_of_json.- Net effect: same columns as before, minus the raw JSON string, plus a new structured
valuecolumn you can dot into.
-- Exploration
SELECT
decoded_key,
value.device AS device, -- field
value.geo.city AS city, -- nested field
value.items AS items, -- array of structs
array_size(value.items) AS number_elements_in_array -- aggregation
FROM bronze_table_struct
ORDER BY number_elements_in_array DESC;
With decoded_value now a STRUCT (value), dot-notation reaches any field directly — value.device for a top-level field, value.geo.city for a nested one — no colon-paths or JSON functions needed anymore. array_size(value.items) counts how many elements are in the items array per row, without unnesting it.
-- Exploding / flattening the array
CREATE OR REPLACE TABLE bronze_explode_array AS
SELECT
decoded_key,
array_size(value.items) AS number_elements_in_array,
explode(value.items) AS item_in_array,
value.items
FROM bronze_table_struct;
explode(collection)is a table-valued generator function: it un-nests anARRAY(orMAP) and returns one row per element, duplicating the other selected columns across those rows.explode(value.items) AS item_in_arrayturns the 3-elementitemsarray shown below into 3 separate rows, each holding one item’s STRUCT.- A
NULLarray produces zero rows, not one row of NULLs — useexplode_outer()instead if rows with a NULL array need to be kept. - The un-exploded
value.itemscolumn is also selected here, so the full 3-element array is visible alongside each individual exploded element — the table below shows the exploded rows.
Result for one order whose items array holds three products (the un-exploded items column is left out here):
| number_elements_in_array | item_in_array.item_id | item_in_array.item_name | item_in_array.price_in_usd |
|---|---|---|---|
| 3 | M_STAN_K |
Standard King Mattress | 1195 |
| 3 | P_FOAM_S |
Standard Foam Pillow | 59 |
| 3 | M_STAN_T |
Standard Twin Mattress | 595 |
Working with a VARIANT column — a bronze-layer pipeline
The VARIANT approach follows the same bronze-layer shape as STRING and STRUCT above, but uses parse_json instead of from_json — and needs no schema at all.
CREATE OR REPLACE TABLE bronze_variant AS
SELECT
decoded_key,
<field_1>,
<field_2>,
parse_json(decoded_value) AS json_variant_value -- conversion into a VARIANT column
FROM bronze_table_decoded;
parse_json(jsonStr)takes a JSON-formatted STRING and returns aVARIANTvalue representing the same data — no schema argument needed, unlikefrom_json. If the string isn’t valid JSON,parse_jsonraises an error (MALFORMED_RECORD_IN_PARSING); usetry_parse_jsoninstead to getNULLback for malformed input rather than failing the query.- This mirrors the STRUCT pipeline exactly — same source table, same decoded column in, one new typed column out — just with a schema-free VARIANT column instead of a fixed STRUCT.
SELECT
json_variant_value,
json_variant_value:device::STRING, -- extract a field & cast it
json_variant_value:items
FROM bronze_variant;
Same colon-path syntax as STRING (and as §10), but on a VARIANT column instead of a STRING one. The difference: a colon-path on VARIANT returns another VARIANT value by default, so the :: shorthand cast (equivalent to CAST(... AS STRING)) is used to get a plain STRING back out — json_variant_value:device::STRING. Left uncast, json_variant_value:items stays a VARIANT holding the nested array.
What does unbase64(col) return, and what extra step gets you back to human-readable text?
unbase64(col) returns BINARY, not a string. Wrapping it in CAST(unbase64(col) AS STRING) converts that binary back into readable text.
Why might a stream's key/value columns arrive base64-encoded in the first place?
Base64 encodes arbitrary bytes as safe, printable ASCII text, avoiding corruption or encoding issues (nulls, special characters) as the data moves through the pipeline — at the cost of the value not being human-readable until it’s decoded.
explode(value.items) AS item_in_array runs against a row whose items array has 3 elements. How many rows does that row become — and what happens if items is NULL instead?
3 rows — one per array element, with the other selected columns duplicated across them. If items is NULL, explode() produces zero rows for that source row (use explode_outer() to keep it as one row of NULLs instead).
On a VARIANT column, what does col:field return by default, and how do you get a plain STRING back?
By default, a colon-path on VARIANT returns another VARIANT value. Append the :: cast shorthand (e.g. col:field::STRING) to cast it to a specific type like STRING.
What’s the key difference between from_json and parse_json?
from_json(jsonStr, schema) requires an explicit schema and returns a typed STRUCT. parse_json(jsonStr) needs no schema and returns a schema-flexible VARIANT value instead.
Sources
- schema_of_json function — Databricks on AWS
- from_json function — Databricks on AWS
- Query variant data — Databricks on AWS
- How is variant different than JSON strings? — Databricks on AWS
- unbase64 function — Databricks on AWS
- explode table-valued generator function — Databricks on AWS
- parse_json function — Databricks on AWS
- Read and write text files — Databricks on AWS
- CONVERT TO DELTA — Databricks on AWS
12. Ingesting Enterprise Data: Managed Connectors & Partner Connect
Objectives 6–9
- Lakeflow Connect Managed Connectors simplify ingestion from enterprise databases and SaaS applications with a fully managed, UI-driven or API-driven experience.
- SaaS ingestion architecture: a serverless Declarative Pipelines job collects credentials from Unity Catalog, reaches out to the public data source, and writes to Streaming Delta Tables.
- Database ingestion architecture: uses an additional Ingestion Gateway (classic compute) to connect to private databases, with staging and state management in a Unity Catalog volume, before a serverless pipeline writes to Streaming Delta Tables.
- Partner Connect provides an alternative when no native managed connector is available, offering a rich ecosystem of partner solutions.
The gap: what CTAS, COPY INTO & Auto Loader don’t cover
| Source | Covered by |
|---|---|
| Cloud object storage | CTAS, COPY INTO, Auto Loader |
| Databases | Gap: needs managed connectors |
| Enterprise (SaaS) applications | Gap: needs managed connectors |
Everything covered so far in this lecture’s code syntax (§8/§9) — CTAS, COPY INTO, and Auto Loader — is built for ingesting from cloud object storage (files in a volume, an S3/ADLS/GCS path). Enterprise databases (on-prem or cloud-hosted MySQL, PostgreSQL, SQL Server, Oracle, etc.) and SaaS applications (Salesforce, Workday, ServiceNow, and similar) are different kinds of sources entirely — they’re not files sitting in object storage, they’re live systems reached over a network connection, each with their own authentication model, API, and change-tracking mechanism. That’s the specific gap Managed Connectors (already introduced conceptually in §2) exist to fill.
Lakeflow Connect Managed Connectors
- Enterprise sources
- Workday
- Salesforce
- ServiceNow
- Google Analytics
- SharePoint
- Dynamics 365
- NetSuite
- PostgreSQL
- SQL Server
- Managed connectors
- Point-and-click UI
- or API
- Data Intelligence Platform
- Managed Connectors are built into Databricks — no separate tool or vendor account needed — and can be set up either through the Databricks UI (“point and click”) or the API/SDK for a code-first workflow.
- Per the official docs, they fall into two broad ingestion architectures depending on the source type: SaaS ingestion (for sources reachable over a public API/endpoint) and Database ingestion (for databases, which need continuous change capture). Both are detailed below.
- Not every connector is Generally Available at the same time — the source material’s own note that “managed connectors in Lakeflow Connect are in various release states” lines up with Databricks’ standard release-stage terminology (Private Preview → Beta → Public Preview → General Availability); worth checking a specific connector’s current release stage before relying on it for production.
Where this happens in the workspace — the Add data UI
Per the official docs, the point-and-click path into a managed connector is: in the Databricks workspace sidebar, click Data Ingestion → this opens the Add data page → under Databricks connectors, pick the source you want (e.g. Salesforce, SQL Server). The Add data page is the single landing spot for every ingestion path covered in this lecture, not just managed connectors — it groups them by category:
- Databricks connectors — the Lakeflow Connect managed connectors covered in this section (SaaS apps, databases).
- File upload / local files / cloud object storage — the CTAS/COPY INTO/Auto Loader-style paths from §3, §8, §9.
- Partner connectors — the Partner Connect ingestion partners covered below.
So the “gap” diagram earlier in this section and the Add data page are really the same map, one conceptual and one literal: cloud storage, databases, and enterprise apps all eventually funnel through this one UI entry point, just into different categories of it.
SaaS ingestion architecture
- Unity Catalog
- Stores the connection credentials
- SaaS source
- Public API, called on each run
- Managed ingestion
- Serverless Declarative Pipeline
- Streaming Delta tables
- A SaaS connector’s pipeline is entirely serverless — a Lakeflow (Serverless) Declarative Pipelines job does the whole job: authenticate using credentials stored in Unity Catalog, call out to the source’s public API/endpoint, and write the result to a Streaming Delta Table.
- Per the docs, SaaS connectors use three components: connections (the Unity Catalog securable object holding auth details), ingestion pipelines, and destination tables — and are automatically compatible with serverless egress controls, since there’s no separate always-on infrastructure to manage.
- No dedicated gateway or staging layer is needed here — unlike the database architecture below — because a public API can simply be called on demand each pipeline run.
Database ingestion architecture
- Database
- On-prem or cloud
- Ingestion gateway
- Classic compute pipeline
- Runs continuously
- Credentials from Unity Catalog
- Staging & state
- Unity Catalog volume
- Purged after 30 days
- Managed ingestion
- Serverless pipeline
- Applies changes (CDC)
- Streaming Delta tables
- Database sources need one extra piece the SaaS architecture doesn’t: the Ingestion Gateway. Per the docs, it runs continuously as its own job on classic compute (not serverless) so it can extract snapshots, change logs, and metadata from the source database in real time — it has to stay running so change logs aren’t truncated at the source before they’re captured.
- The gateway writes what it captures to staging storage — a Unity Catalog volume that temporarily holds the extracted data. This decouples continuous change capture (gateway) from the actual load into Delta tables (pipeline), so the pipeline can run on its own schedule independently of the gateway’s constant capture. Per the docs, staged data auto-purges after 30 days.
- From there, a separate Managed Ingestion job — serverless, same as the SaaS path — reads the staged data and writes/merges it into the destination Streaming Delta Table(s), using CDC to apply just the changes rather than reprocessing everything.
- Officially documented managed database connectors (as of this writing) cover MySQL, PostgreSQL, Microsoft SQL Server, and Oracle, each with a CDC mechanism specific to that database engine.
Partner Connect — the alternative
Partner Connect lets you create trial accounts with select Databricks technology partners and connect your Databricks workspace to partner solutions directly from the Databricks UI — letting you try a partner’s ingestion (or BI, or other) tooling against your own lakehouse data before committing to full adoption. It’s the fallback for sources that don’t have a native Lakeflow Connect managed connector yet.
Why can't CTAS, COPY INTO, or Auto Loader ingest directly from a Salesforce account the same way they ingest from an S3 bucket?
Those three methods are built for cloud object storage (files at a path). A SaaS application like Salesforce is a live system reached over its own API with its own authentication model — not a file source — which is exactly the gap Managed Connectors (SaaS ingestion architecture) are built to handle.
In the Databricks workspace UI, where do you go to start setting up a Salesforce managed connector, and what else lives on that same page?
Sidebar → Data Ingestion → the Add data page → under Databricks connectors, pick Salesforce. The same Add data page also holds the file-upload/cloud-storage ingestion paths and Partner Connect’s ingestion partners — it’s the single landing point for every ingestion method covered in this lecture.
What's the one architectural component a database connector needs that a SaaS connector doesn't, and why?
The Ingestion Gateway. It has to run continuously on classic compute to capture the database’s change logs/snapshots before they’re truncated at the source — a SaaS connector can just call a public API on demand each run, with no continuous capture needed.
In the database ingestion architecture, what's the role of the Unity Catalog volume, and how long is data kept there?
It’s the staging layer: the Ingestion Gateway writes extracted snapshots/change logs there, decoupling continuous capture from the actual (serverless, scheduled) load into Delta tables. Staged data automatically purges after 30 days.
You need to ingest from a source with no native Lakeflow Connect managed connector. What's the built-in alternative, and what tier/permissions does it require?
Partner Connect — it lets you create a trial account with a Databricks technology partner and connect it to your workspace from the Databricks UI. Requires a Premium tier (or higher) account, and workspace admin privileges to set up a new connection.
Sources
- What is Lakeflow Connect? — Databricks on AWS
- Choose a standard connector — Databricks on AWS
- Lakeflow Connect connector concepts — Databricks on AWS
- Managed database connectors — Databricks on AWS
- Ingest data from Salesforce (Add data UI steps) — Databricks on AWS
- What is Databricks Partner Connect? — Databricks on AWS
- Connect to ingestion partners using Partner Connect — Databricks on AWS
- Technology partners — Databricks on AWS
- Release types — Databricks on AWS
13. Other Data Integration & Sharing Features
Objectives 10–11
Beyond the ingestion methods already covered (§2–§4, §8, §9, §12), Databricks offers a handful of additional features for integrating with and sharing data: Lakehouse Federation (query external sources without moving the data), Zerobus (a direct, high-throughput event-ingestion API), Delta Sharing (secure read-only data sharing), and the Databricks Marketplace (an open exchange for data products, built on top of that sharing protocol).
Lakehouse Federation
- Lakehouse Federation is Databricks’ query federation platform — it gives governed, read-only access to external data through Unity Catalog foreign catalogs, without copying or moving that data into Databricks first. This is the “data virtualization” framing from the lecture: query it where it lives.
- Queries are pushed down to the external system where possible (via JDBC for relational sources) — so filtering/aggregation happens at the source, and only the result set comes back to Databricks — with Unity Catalog governance (permissions, auditing) still applying on top.
- Supported external sources include MySQL, PostgreSQL, SQL Server, Oracle, Teradata, Amazon Redshift, Snowflake, Google BigQuery, Azure Synapse, Salesforce Data 360, and other Databricks workspaces/instances — plus catalog-level federation for Hive metastore, AWS Glue, and Palantir Foundry.
- Best suited for ad hoc reporting, BI, and proof-of-concept access to operational databases — not a replacement for actually ingesting data you’ll query heavily or repeatedly (the ingestion methods in §2–§4/§12 are still the better fit there).
Zerobus
- Zerobus Ingest is a Lakeflow Connect component for writing event data directly into Unity Catalog Delta tables at high, sustained throughput — skipping any intermediate queuing/streaming system (like standing up your own Kafka cluster) between the producer and the lakehouse.
- It offers multiple interfaces depending on the producer: gRPC-based SDKs for the highest sustained throughput (high-volume streaming producers), a stateless REST API (suited to large fleets of lightweight/“chatty” edge devices), OpenTelemetry support (for systems already emitting traces/logs/metrics), and a Kafka-compatible API (in beta, for migrating existing Kafka producers).
- Supports JSON, Protocol Buffers, and Apache Arrow message formats, with schema validation and durable fallback mechanisms built in.
Delta Sharing → OpenSharing
Delta Sharing is the open protocol for secure, read-only data sharing — it lets you share tables, views, volumes, notebooks, and models with recipients outside your organization, regardless of whether they use Databricks, without copying the data into their system. Real-time: when the provider updates the shared data, recipients see the update in near real time.
Databricks Marketplace
Databricks Marketplace is an open exchange where data providers, software vendors, and technology partners publish offerings that customers can discover, evaluate, and connect to directly from their own workspace — browsing and consuming data products without leaving Databricks. It uses OpenSharing (Delta Sharing) under the hood for the actual secure data exchange, and integrates with Partner Connect for pre-configured connections to select technology partners.
You need to run a one-off BI query against a live PostgreSQL database without copying its data into Databricks. What feature fits, and what actually executes the filtering?
Lakehouse Federation (query federation) — it accesses the external database through a Unity Catalog foreign catalog, and pushes the query (filtering/aggregation) down to run in the external database itself via JDBC, returning only the result set.
An IoT fleet of edge devices needs to write event data directly into a Delta table at high throughput, with minimal infrastructure in between. What's the fit?
Zerobus Ingest — it writes events directly into Unity Catalog Delta tables at high sustained throughput via gRPC SDKs, REST API, OpenTelemetry, or a (beta) Kafka-compatible API, without needing a separate queuing/streaming system.
What's the current Databricks product name for older training material calls "Delta Sharing," and is Delta Sharing itself gone?
The current product/feature name is OpenSharing. Delta Sharing isn’t gone — it’s the underlying open protocol OpenSharing is built on and expands (from data sharing to the full AI asset stack), and it’s still the term other platforms use for their implementation of that same protocol.
What's the relationship between Databricks Marketplace and Delta Sharing/OpenSharing?
Marketplace is the open exchange/catalog UI where data providers list products (datasets, models, notebooks, apps, MCP servers) for customers to discover and connect to — it uses OpenSharing (Delta Sharing) as the underlying protocol for the actual secure data exchange.
Sources
- Connect to external databases and catalogs (Lakehouse Federation) — Databricks on AWS
- What is query federation? — Databricks on AWS
- Use the Zerobus Ingest connector — Databricks on AWS
- OpenSharing — Databricks on AWS
- Introducing OpenSharing: the Next Evolution of Delta Sharing — Databricks Blog
- Databricks Marketplace — Databricks on AWS
Exam Objective Crosswalk
Verbatim bullets from the official exam guide (May 4, 2026 version) that this lecture’s content supports.
| Exam Guide Objective (verbatim) | Section | Lecture Topic |
|---|---|---|
| Understand the core components of the Databricks Data Intelligence Platform, such as its architecture, Delta Lake, and Unity Catalog. | §1 (6%) | §1, §5, §6, §7 |
| Enable and detail data ingestion patterns, including batch, streaming, and incremental loading, and import data from sources such as local files, Lakeflow Connect standard connectors, and Lakeflow Connect managed connectors. | §2 (21%) | §1–§4 |
| Configure Lakeflow Connect to reliably ingest data from diverse enterprise sources into Unity Catalog–governed tables. | §2 (21%) | §1, §2, §12 |
| Prioritize between Auto Loader, Lakeflow Connect (standard and managed connectors), partner connectors, and other ingestion methods based on technical requirements such as data volume, ingestion frequency, data types, and governance needs with Unity Catalog. | §2 (21%) | §4, §12 |
| Ingest semi-structured and unstructured data (for example, JSON and nested data) via Lakeflow Connect and other managed connectors into Unity Catalog–governed Delta tables. | §2 (21%) | §2, §3, §11 |
Practice Questions
Sample multiple-choice questions covering §2, §8, §12. Answers verified against official docs where checkable; see notes under each.
You are designing an ingestion pipeline using Lakeflow Connect in Databricks.
- Some of your data already exists as files in cloud object storage that Databricks can directly access.
- Other data lives in an external system, such as a transactional database or SaaS application, and has not yet been landed in cloud storage.
Which choice correctly matches each scenario to the connector type you should use?
- Use managed connectors for data already in cloud storage, and standard connectors for data in external systems
- Use standard connectors for data already in cloud storage, and managed connectors for data in external systems
- Use standard connectors for both data in cloud storage and external systems
- Use managed connectors for both data in cloud storage and external systems
Show answer
Correct: B. Standard connectors are built for cloud object storage/message buses (§2) — the data’s already reachable as files. Managed connectors are built for external systems (databases, SaaS apps) that haven’t landed in storage yet (§2, §12) — they handle source-specific auth, CDC, and schema evolution automatically. Matches the notes exactly.
You need to ingest CSV files that have the following characteristics:
- Files are semicolon (
;) delimited instead of comma delimited - Files contain headers in the first row
- Files may contain malformed data that should be captured for later analysis
- You want to enforce a specific schema:
order_id BIGINT, customer_email STRING, order_total DECIMAL(10,2)
Which read_files() configuration correctly handles all these requirements?
SELECT * FROM read_files( "/path/to/files", format => "csv", sep => ";", header => true, schema => "order_id:BIGINT, customer_email:STRING, order_total:DECIMAL(10,2)", rescuedDataColumn => "_rescued_data" );SELECT * FROM read_files( "/path/to/files", format => "csv", separator => ";", headers => true, enforceSchema => "order_id BIGINT, customer_email STRING, order_total DECIMAL(10,2)", rescueColumn => "_rescued_data" );SELECT * FROM read_files( "/path/to/files", format => "csv", delimiter => ";", header => true, schema => "order_id BIGINT, customer_email STRING, order_total DECIMAL(10,2)" );SELECT * FROM read_files( "/path/to/files", format => "csv", sep => ";", header => true, schema => "order_id BIGINT, customer_email STRING, order_total DECIMAL(10,2)", rescuedDataColumn => "_rescued_data" );
Show answer
Correct: D. Verified option-by-option against the official CSV DataFrameReader/read_files reference: the delimiter option is sep (not separator or delimiter — both A and D use sep correctly, C and B use the wrong name), header (not headers) is correct, the schema value must be a space-separated DDL string like "col TYPE, col TYPE" — not colon-separated (col:TYPE), which rules out A — and the rescued-data option is spelled rescuedDataColumn specifically (not rescueColumn or enforceSchema, ruling out B), taking a string value for the desired column name. C is missing the rescued-data option entirely, which the scenario explicitly calls for (“malformed data that should be captured”). Only D gets every option name and format right.
What is the purpose of the Ingestion Gateway component in the Database ingestion flow?
- Storing final Streaming Tables
- Managing user permissions
- Connecting to the source database
- Running BI reports
Show answer
Correct: C. Per §12: the Ingestion Gateway runs continuously on classic compute specifically to connect to the source database and extract snapshots/change logs/metadata before they’re truncated at the source. It doesn’t store the final Streaming Delta Tables (that’s the separate, serverless Managed Ingestion step), doesn’t manage permissions (that’s Unity Catalog), and has nothing to do with BI reporting.
How is Partner Connect commonly used when ingesting data into Databricks?
- To write custom Spark code for every external data source
- To manually upload files from a local machine into Unity Catalog volumes
- To configure and launch partner ingestion tools that load data from external systems into Databricks
- To replace built-in Databricks ingestion features like Auto Loader
Show answer
Correct: C. Per §12: Partner Connect lets you create trial accounts with technology partners and connect your workspace to their solutions from the Databricks UI — it’s the fallback for sources with no native managed connector, not a code-writing tool, not manual file upload, and not a replacement for Auto Loader (it’s an alternative for sources Auto Loader/managed connectors don’t cover).
What is the purpose of the rescued data column when ingesting data into Databricks?
- To log ingestion errors
- To handle records that don’t match the schema of the target table
- To store duplicate records
- To store metadata about ingestion jobs
Show answer
Correct: B. Per §8/§10: _rescued_data captures columns/values that don’t match the expected schema (missing fields, type mismatches, case mismatches) as a JSON-formatted string, so mismatched data is preserved instead of silently dropped — it’s not an error log, not a dedup mechanism, and not job metadata (that’s _metadata, Question 6).
When ingesting data into a Bronze table using the _metadata column, which of the following metadata information can be extracted from input files?
- File content and data schema information
- File size and file permissions only
- File name, file modification time, and file path
- Only the file creation timestamp
Show answer
Best available answer: C. A and B are wrong outright — _metadata doesn’t expose file content/schema, and there’s no “file permissions” field at all. D is wrong too — there’s no file creation timestamp field, only modification time. C is correct as far as it goes, but worth knowing it’s incomplete: per the official docs, the full _metadata struct also includes file_size, file_block_start, and file_block_length alongside file_path, file_name, and file_modification_time — so C names three real fields correctly but isn’t the complete list. Worth flagging if this question’s phrasing (“which of the following can be extracted”) shows up expecting a single best choice among imperfect options, as it does here.
You are using MERGE INTO in Databricks SQL to upsert data from a source table into a Delta target table. The source table may occasionally include new columns that do not yet exist in the target table. You want the merge operation to handle these changes automatically.
Which statement best accomplishes this?
MERGE INTO target t USING source s ON t.id = s.id WHEN MATCHED THEN UPDATE SET * WHEN NOT MATCHED THEN INSERT *; ALTER TABLE target ADD COLUMNS (...);MERGE INTO target t USING source s ON t.id = s.id WHEN MATCHED THEN UPDATE SET * WHEN NOT MATCHED THEN INSERT *;MERGE INTO target t USING source s ON t.id = s.id WHEN MATCHED THEN UPDATE SET t.col1 = s.col1;MERGE WITH SCHEMA EVOLUTION INTO target t USING source s ON t.id = s.id WHEN MATCHED THEN UPDATE SET * WHEN NOT MATCHED THEN INSERT *;
Show answer
Correct: D. Per the official MERGE INTO syntax (§8): the full statement grammar is MERGE [WITH SCHEMA EVOLUTION] INTO target USING source ON ... — the optional WITH SCHEMA EVOLUTION clause is exactly what tells the merge to automatically add new source columns to the target schema during the operation. A bolts on a separate manual ALTER TABLE after the fact (works, but isn’t “automatic,” and you’d need to know the new columns in advance to write it). B is plain MERGE INTO with no schema evolution — new source columns would be dropped or cause a schema mismatch, not automatically added. C explicitly lists only one column, the opposite of automatic handling.
Independent study notes, not affiliated with or endorsed by Databricks. Product names belong to their owners. Databricks changes quickly, so check anything important against the official documentation.