How Semantica Ingests Data from Databricks and Snowflake: A Complete Technical Guide
Semantica ingests data from Databricks and Snowflake through dedicated connector modules that wrap native client libraries, stream query results through cursors, and convert tabular data into JSON-serializable documents ready for downstream processing.
Semantica provides purpose-built ingestion pipelines for enterprise data warehouses. Whether you're pulling from Databricks Unity Catalog with Delta Lake or a Snowflake warehouse, the framework follows a consistent two-layer architecture—connectors handle authentication and transport, while ingestors manage extraction, metadata discovery, and document transformation.
Databricks Data Ingestion Architecture
The Databricks ingestion pipeline centers on DatabricksConnector and DatabricksIngestor in semantica/ingest/databricks_ingestor.py. This module supports both SQL warehouse connections and Unity Catalog metadata operations.
DatabricksConnector: Authentication and Connection Management
DatabricksConnector wraps two native clients:
- databricks-sql-connector for SQL execution
- databricks-sdk for Unity Catalog operations
The connector supports multiple authentication flows:
| Method | Use Case |
|---|---|
| Personal access token | Quick scripts, development |
| OAuth M2M (service principal) | Production automation |
The connector lazily instantiates a WorkspaceClient for catalog operations and reuses existing SQL connections to minimize overhead.
DatabricksIngestor: Data Extraction and Metadata Discovery
DatabricksIngestor exposes two primary extraction methods:
ingest_table(table_name, limit=None)— Full table scan with optional row limitingest_query(sql, limit=None)— Arbitrary SQL execution
Both methods:
- Build safe, escaped SQL statements
- Stream results through a cursor (memory-efficient)
- Convert rows to JSON-serializable dictionaries
Additional metadata helpers expose Unity Catalog information:
list_catalogs()— All accessible catalogslist_schemas(catalog)— Schemas within a cataloglist_tables(catalog, schema)— Tables in a schemaget_table_schema(table)— Column definitions and constraintsget_table_lineage(table)— Data lineage from Unity Catalog
Practical Databricks Example
from semantica.ingest import DatabricksIngestor
# Initialize with personal access token
ingestor = DatabricksIngestor(
host="https://adb-123.azuredatabricks.net",
token="dapi-xxxxxxxxxxxx",
http_path="/sql/1.0/warehouses/warehouse-abc",
catalog="main",
schema="default",
)
# Ingest table data
data = ingestor.ingest_table("customers", limit=5000)
# Convert to documents for downstream processing
documents = ingestor.export_as_documents(
data,
text_fields=["first_name", "last_name", "email"]
)
# Explore catalog metadata
catalogs = ingestor.list_catalogs()
schemas = ingestor.list_schemas(catalog="main")
tables = ingestor.list_tables(catalog="main", schema="default")
Snowflake Data Ingestion Architecture
The Snowflake pipeline in semantica/ingest/snowflake_ingestor.py mirrors the Databricks design with SnowflakeConnector and SnowflakeIngestor.
SnowflakeConnector: Flexible Authentication
SnowflakeConnector wraps snowflake-connector-python and accepts connection parameters via arguments or environment variables. Supported authentication modes:
- Password-based
- Key-pair (encrypted or unencrypted)
- OAuth token
- External browser SSO
SnowflakeIngestor: Extraction with Validation
SnowflakeIngestor provides the same core API as its Databricks counterpart:
ingest_table(table, where=None, order_by=None, limit=None, offset=None)ingest_query(sql, limit=None)
Security features include:
- Double-quote identifier escaping
WHEREclause validation (prevents injection)ORDER BYclause validation- Optional pagination via
LIMIT/OFFSET
Metadata operations:
list_tables(database, schema)get_table_schema(table)— Returns column definitions and primary keys
Practical Snowflake Example
from semantica.ingest import SnowflakeIngestor
# OAuth authentication
ingestor = SnowflakeIngestor(
account="myaccount",
user="myuser",
authenticator="oauth",
token="my-oauth-token",
warehouse="COMPUTE_WH",
database="SALES_DB",
schema="PUBLIC",
)
# Ingest with filter and limit
data = ingestor.ingest_table(
"orders",
where="order_date > '2024-01-01'",
limit=10000
)
# Export to document format
documents = ingestor.export_as_documents(
data,
text_fields=["order_id", "description"]
)
# Inspect schema metadata
schema_info = ingestor.get_table_schema("CUSTOMERS")
print(schema_info["columns"])
print("Primary keys:", schema_info["primary_keys"])
Common Infrastructure: Progress Tracking and Logging
Both ingestion modules share supporting infrastructure from semantica/utils/:
| Component | File | Purpose |
|---|---|---|
| Progress tracking | semantica/utils/progress_tracker.py |
get_progress_tracker() emits start/stop events for long-running operations |
| Logging | semantica/utils/logging.py |
get_logger() provides structured debug, info, and error messages |
| Exceptions | semantica/utils/exceptions.py |
ProcessingError and ValidationError for consistent error handling |
Data Representation and Document Export
Ingested data flows into simple @dataclass containers:
DatabricksData—rows,columns,row_count,metadata(source query, catalog, schema)SnowflakeData— Same structure, Snowflake-specific metadata
Both ingestors implement export_as_documents(data, text_fields), which transforms tabular results into a list of document dictionaries:
{
"id": "row_0",
"metadata": {"first_name": "Alice", "last_name": "Smith", ...},
"text": "Alice Smith alice@example.com"
}
This uniform format feeds directly into Semantica's knowledge-graph processing pipelines.
Summary
- Two-layer architecture: Connectors handle transport/authentication; ingestors manage extraction and transformation
- Databricks ingestion (
semantica/ingest/databricks_ingestor.py): Unity Catalog + Delta Lake support with lineage tracking - Snowflake ingestion (
semantica/ingest/snowflake_ingestor.py): Full warehouse connectivity with SQL validation - Shared infrastructure: Progress tracking and logging from
semantica/utils/ - Unified output:
DatabricksDataandSnowflakeDataconvert to documents viaexport_as_documents()
Frequently Asked Questions
What authentication methods does Semantica support for Databricks?
Semantica's DatabricksConnector supports personal access tokens and OAuth M2M (service principal) authentication according to the semantica-agi/semantica source code. Personal access tokens suit development workflows, while OAuth M2M enables secure production automation without long-lived credentials.
How does Semantica handle large tables without memory issues?
Both DatabricksIngestor and SnowflakeIngestor stream query results through database cursors rather than loading entire tables into memory. The ingest_table() and ingest_query() methods process rows iteratively, and you can apply limit and offset parameters for controlled pagination.
Can I extract metadata and lineage from Databricks Unity Catalog?
Yes. DatabricksIngestor provides dedicated methods: list_catalogs(), list_schemas(), list_tables(), and get_table_lineage() expose Unity Catalog's full metadata and lineage graph. These operations use a lazily-initialized WorkspaceClient from the databricks-sdk.
What security validations does the Snowflake ingestor apply?
SnowflakeIngestor escapes all identifiers with double quotes and validates WHERE and ORDER BY clauses against injection patterns. The connector itself supports password, key-pair, OAuth, and external browser SSO authentication for flexible security postures.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →