Skip to content

Document Processing Pipeline

Skill: databricks-ai-functions

A multi-stage document processing pipeline that parses PDFs and images, classifies document types, extracts structured fields, and matches entities against master data — all using SQL AI Functions inside a Spark Declarative Pipeline. Each stage uses the cheapest function that can handle the job: ai_parse_document for OCR, ai_classify v2 for routing, ai_extract v2 for flat or nested fields, and ai_query only when the output schema exceeds ai_extract limits (128 fields, 7 nesting levels).

“Write a SDP pipeline in Python that parses documents from a landing volume, classifies them with ai_classify v2, extracts invoice fields with ai_extract v2 typed schemas, and routes errors to a sidecar table.”

import dlt
import yaml
from pyspark.sql.functions import expr, col, from_json
CFG = yaml.safe_load(open("/Workspace/path/to/config.yml"))
ENDPOINT = CFG["models"]["default"]
PROMPT = CFG["prompts"]["extract_invoice"]
# Stage 1: Parse binary documents
@dlt.table(comment="Parsed document text from landing volume")
def raw_parsed():
return (
spark.read.format("binaryFile").load(CFG["volumes"]["input"])
.withColumn("parsed", expr("ai_parse_document(content, MAP('version', '2.0'))"))
.withColumn("text_blocks", expr("""
concat_ws('\n', transform(
parsed:document:elements,
e -> e:content::STRING
))
"""))
.selectExpr(
"path",
"text_blocks",
"parsed:error_status AS parse_error",
)
.filter("parse_error IS NULL")
)
# Stage 2: Classify document type with ai_classify v2
@dlt.table(comment="Document type classification")
def classified_docs():
return (
dlt.read("raw_parsed")
.withColumn(
"doc_type",
expr("""
ai_classify(
text_blocks,
'["invoice", "purchase_order", "receipt", "contract", "other"]',
MAP('version', '2.0')
):response[0]::STRING
""")
)
)
# Stage 3a: Flat field extraction with ai_extract v2 typed schema
@dlt.table(comment="Flat header fields extracted from invoices")
def extracted_flat():
return (
dlt.read("classified_docs")
.filter("doc_type = 'invoice'")
.filter("text_blocks IS NOT NULL")
.withColumn(
"result",
expr("""
ai_extract(
text_blocks,
'{
"invoice_number": {"type": "string"},
"vendor_name": {"type": "string"},
"issue_date": {"type": "string", "description": "dd/mm/yyyy"},
"total_amount": {"type": "number"},
"tax_id": {"type": "string"}
}',
MAP('version', '2.0')
)
""")
)
.selectExpr(
"path", "doc_type", "text_blocks",
"result:response AS header",
"result:error_message::STRING AS extract_error"
)
)
# Stage 3b: Nested line items via ai_query (last resort)
@dlt.table(comment="Nested line items -- ai_query used only for the array schema")
def extracted_line_items():
return (
dlt.read("extracted_flat")
.filter("extract_error IS NULL")
.withColumn(
"ai_response",
expr(f"""
ai_query(
'{ENDPOINT}',
concat('{PROMPT.strip()}', '\\n\\nDocument text:\\n', LEFT(text_blocks, 6000)),
responseFormat => '{{"type":"json_object"}}',
failOnError => false
)
""")
)
.withColumn(
"line_items",
from_json(
col("ai_response.response"),
"STRUCT<line_items:ARRAY<STRUCT<item_code:STRING, description:STRING, "
"quantity:DOUBLE, unit_price:DOUBLE, total:DOUBLE>>>"
)
)
.select("path", "doc_type", "header", "line_items",
col("ai_response.error").alias("extraction_error"))
)
# Stage 4: Success output
@dlt.table(comment="Processed documents ready for downstream consumption")
def processed_docs():
return dlt.read("extracted_line_items").filter("extraction_error IS NULL")
# Stage 5: Error sidecar
@dlt.table(comment="Failed extractions for review and reprocess")
def processing_errors():
return (
dlt.read("extracted_flat")
.filter("extract_error IS NOT NULL")
.select("path", "doc_type", col("extract_error").alias("error"))
.unionByName(
dlt.read("extracted_line_items")
.filter("extraction_error IS NOT NULL")
.select("path", "doc_type", col("extraction_error").alias("error"))
)
)

Key decisions:

  • ai_parse_document with MAP('version', '2.0') — the v2 schema returns elements under parsed:document:elements and surfaces errors at parsed:error_status. Always use v2 for new pipelines; the legacy parsed:pages[*].elements[*] shape is going away.
  • ai_classify v2 with JSON-string labels — pass '["invoice", ...]' plus MAP('version', '2.0'), then extract the top label with :response[0]::STRING. The legacy array(...) form is deprecated; v2 returns VARIANT with response and error_message fields.
  • ai_extract v2 with typed schema — declare types (string, number, enum) so downstream tables get clean STRUCTs instead of raw VARIANT. Up to 128 fields and 7 nesting levels, which covers most invoice headers without falling back to ai_query.
  • ai_query only when schema exceeds ai_extract limits — nested arrays like line_items blow past 7 nesting levels. That is the one case where ai_query with responseFormat is the right tool.
  • failOnError => false everywhere — a single malformed document should not kill the entire batch. Route errors to a sidecar table for retry.
  • LEFT(text_blocks, 6000) caps input length — long documents silently truncate or exceed context windows. Truncate explicitly.

“Write SQL to parse documents from a volume, explode elements into individual chunks, and store them in a table for downstream processing.”

CREATE OR REPLACE TABLE catalog.schema.parsed_chunks AS
WITH parsed AS (
SELECT
path,
ai_parse_document(content) AS doc
FROM read_files('/Volumes/catalog/schema/volume/docs/', format => 'binaryFile')
),
elements AS (
SELECT
path,
explode(variant_get(doc, '$.document.elements', 'ARRAY<VARIANT>')) AS element
FROM parsed
)
SELECT
md5(concat(path, variant_get(element, '$.content', 'STRING'))) AS chunk_id,
path AS source_path,
variant_get(element, '$.content', 'STRING') AS content,
variant_get(element, '$.type', 'STRING') AS element_type,
current_timestamp() AS parsed_at
FROM elements
WHERE variant_get(element, '$.content', 'STRING') IS NOT NULL
AND length(trim(variant_get(element, '$.content', 'STRING'))) > 10;

This pure-SQL approach is useful when you want to chunk documents for RAG or vector search without a full SDP pipeline. The length > 10 filter drops whitespace-only elements that add noise.

Streaming ingestion with incremental parsing

Section titled “Streaming ingestion with incremental parsing”

“Write PySpark to set up a streaming job that parses new documents as they land in a volume, using v2 schema and exactly-once checkpoints.”

from pyspark.sql.functions import col, current_timestamp, expr
files_df = (
spark.readStream.format("binaryFile")
.option("pathGlobFilter", "*.{pdf,jpg,jpeg,png}")
.option("recursiveFileLookup", "true")
.load("/Volumes/catalog/schema/volume/docs/")
)
parsed_df = (
files_df
.repartition(8, expr("crc32(path) % 8"))
.withColumn("parsed", expr("""
ai_parse_document(content, map(
'version', '2.0',
'descriptionElementTypes', '*'
))
"""))
.withColumn("parsed_at", current_timestamp())
.select("path", "parsed", "parsed_at")
)
(
parsed_df.writeStream.format("delta")
.outputMode("append")
.option("checkpointLocation", "/Volumes/catalog/schema/checkpoints/01_parse")
.option("mergeSchema", "true")
.trigger(availableNow=True)
.toTable("catalog.schema.parsed_documents_raw")
)

repartition by file hash parallelizes ai_parse_document across workers; without it, a single partition processes files sequentially. trigger(availableNow=True) processes all pending files then stops — ideal for scheduled batch jobs. The checkpoint guarantees exactly-once processing so re-runs do not re-parse unchanged documents. For a production-ready two-stage variant that separates parse from text extraction with independent checkpoints, see the bundle example.

“Write SQL to match extracted vendor names against a master data table using fuzzy similarity scoring.”

SELECT
e.path,
e.invoice.vendor_name AS extracted_vendor,
m.vendor_id,
m.vendor_name AS master_vendor,
ai_similarity(e.invoice.vendor_name, m.vendor_name) AS match_score
FROM catalog.schema.processed_docs e
CROSS JOIN catalog.schema.vendor_master m
WHERE ai_similarity(e.invoice.vendor_name, m.vendor_name) > 0.80
ORDER BY match_score DESC;

ai_similarity returns a 0-1 score based on semantic similarity. A threshold of 0.80 catches common variations (“Acme Corp” vs “ACME Corporation”) while filtering noise. For large vendor tables, pre-filter with a LIKE clause before running similarity to reduce cross-join cardinality.

  • Passing raw binary to ai_query produces garbage — always parse documents with ai_parse_document first, then feed the extracted text to ai_query or ai_extract. Binary content is not text.
  • Skipping MAP('version', '2.0') on v2 functions — ai_parse_document, ai_classify, and ai_extract all gained v2 schemas with a different response shape (:response, :error_message, parsed:document:elements). Pass MAP('version', '2.0') and read v2 paths; legacy callers still work but will eventually break.
  • ai_extract only when fields exceed 128 or nesting exceeds 7 levels — otherwise reach for ai_extract v2 with a typed schema. Falling back to ai_query for flat fields wastes endpoint quota and produces less consistent output than the optimized task-specific endpoint.
  • explode on a VARIANT — explode() requires ARRAY. Use variant_get(doc, '$.document.elements', 'ARRAY<VARIANT>') to cast before exploding when chunking for RAG.
  • ai_parse_document is the bottleneck — it is the slowest stage in the pipeline. Repartition before calling it to parallelize across executors, and use trigger(availableNow=True) for batch scheduling.
  • Prompt drift across config changes — externalizing prompts to a config.yml makes them versionable and testable. Hardcoded prompts scattered across pipeline stages are hard to audit and easy to break.