Document Processing Pipeline
Skill: databricks-ai-functions
What You Can Build
Section titled “What You Can Build”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).
In Action
Section titled “In Action”“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 dltimport yamlfrom 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_documentwithMAP('version', '2.0')— the v2 schema returns elements underparsed:document:elementsand surfaces errors atparsed:error_status. Always use v2 for new pipelines; the legacyparsed:pages[*].elements[*]shape is going away.ai_classifyv2 with JSON-string labels — pass'["invoice", ...]'plusMAP('version', '2.0'), then extract the top label with:response[0]::STRING. The legacyarray(...)form is deprecated; v2 returns VARIANT withresponseanderror_messagefields.ai_extractv2 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 toai_query.ai_queryonly when schema exceedsai_extractlimits — nested arrays likeline_itemsblow past 7 nesting levels. That is the one case whereai_querywithresponseFormatis the right tool.failOnError => falseeverywhere — 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.
More Patterns
Section titled “More Patterns”Parse and chunk documents in pure SQL
Section titled “Parse and chunk documents in pure SQL”“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 ASWITH 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_atFROM elementsWHERE 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.
Fuzzy vendor matching with ai_similarity
Section titled “Fuzzy vendor matching with ai_similarity”“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_scoreFROM catalog.schema.processed_docs eCROSS JOIN catalog.schema.vendor_master mWHERE ai_similarity(e.invoice.vendor_name, m.vendor_name) > 0.80ORDER 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.
Watch Out For
Section titled “Watch Out For”- Passing raw binary to
ai_queryproduces garbage — always parse documents withai_parse_documentfirst, then feed the extracted text toai_queryorai_extract. Binary content is not text. - Skipping
MAP('version', '2.0')on v2 functions —ai_parse_document,ai_classify, andai_extractall gained v2 schemas with a different response shape (:response,:error_message,parsed:document:elements). PassMAP('version', '2.0')and read v2 paths; legacy callers still work but will eventually break. ai_extractonly when fields exceed 128 or nesting exceeds 7 levels — otherwise reach forai_extractv2 with a typed schema. Falling back toai_queryfor flat fields wastes endpoint quota and produces less consistent output than the optimized task-specific endpoint.explodeon a VARIANT —explode()requiresARRAY. Usevariant_get(doc, '$.document.elements', 'ARRAY<VARIANT>')to cast before exploding when chunking for RAG.ai_parse_documentis the bottleneck — it is the slowest stage in the pipeline. Repartition before calling it to parallelize across executors, and usetrigger(availableNow=True)for batch scheduling.- Prompt drift across config changes — externalizing prompts to a
config.ymlmakes them versionable and testable. Hardcoded prompts scattered across pipeline stages are hard to audit and easy to break.