Data Service — Schema & Fetch Dataflow
Scenario A — Single-Table CREATE
POST /api/apps/{slug}/datatables/
HTTP Request
│
▼
AppDataTableViewSet.create() views.py
│
▼
DataTableSerializer.validate(attrs) serializers.py
├─ normalize_frictionless_fk_actions(schema)
├─ adapter.normalize_schema_for_validation(schema) vector → any + x-vector-metadata
├─ adapter.validate(normalized) Frictionless type check
├─ SqlMapper.enrich_search_fields() if search_fields present
└─ SqlMapper.enrich_vector_fields() injects HNSW index entry
│ stores canonical any+meta form as attrs["json_schema"]
▼
DataTableSerializer.create(validated_data)
├─ build_physical_table_name(app.slug, name)
└─ super().create() → DataTable.save()
│
▼
DataTable.save() models.py
├─ schema_hash = sha256(json_schema)
├─ physical_table_name = build_physical_table_name()
├─ full_clean() → clean() → adapter.validate()
└─ super().save() → DB INSERT
│
├─ [on_commit] post_save signal
│ ├─ SchemaCache.invalidate()
│ └─ SATableBuilder.clear_cache(schema_name, table_name)
│
└─ DataTable.materialize()
├─ ForeignKeyResolutionService.normalize_fk_references()
├─ FrictionlessAdapter.to_sqlalchemy_table(schema)
│ ├─ swap tsvector/search_vector → string for from_descriptor()
│ ├─ FrictionlessSchema.from_descriptor(parse_schema)
│ └─ SqlMapper.write_schema(fs, _schema_dict=schema)
│ ├─ write_type(any, descriptor) → PgVector / HALFVEC
│ └─ _add_custom_indexes() → HNSW index object
├─ _ensure_pgvector_extension()
│ └─ CREATE EXTENSION IF NOT EXISTS vector SCHEMA public
└─ schema_manager.create_table(sa_table)
└─ Alembic CreateTableOp → DDL
CREATE TABLE ... (embedding vector(384), ...)
CREATE INDEX ... USING hnsw (embedding vector_cosine_ops)
WITH (m=16, ef_construction=64)
Scenario B — Bulk/Package CREATE
POST /api/apps/{slug}/datatables/create_schema/
HTTP Request
│
▼
AppDataTableViewSet.create_schema() views.py
│
▼
DataPackageSerializer.is_valid()
└─ DataPackageResourceSerializer.validate(attrs) per resource
├─ normalize_frictionless_fk_actions(schema)
├─ adapter.normalize_schema_for_validation() vector → any ← line ~535
├─ SqlMapper.enrich_hierarchy() / enrich_graph() if hierarchy
├─ SqlMapper.enrich_search_fields() if search_fields
├─ SqlMapper.enrich_vector_fields() HNSW index entry
└─ adapter.validate(schema)
│ attrs["schema"] is canonical any+meta
▼
DataPackageSerializer.create_datatables(app, user)
└─ DataPackageService.create_tables_bulk(resources)
├─ logical → physical name resolution + UUID PK normalization
├─ FK dependency ordering (Alembic sorted_tables)
└─ per table:
├─ [new] DataTable.objects.create() → DataTable.save() → materialize()
└─ [existing] datatable.save() → apply_migration()
└─ Alembic diff → ADD / DROP / ALTER COLUMN + index ops
Scenario C — Schema UPDATE
PATCH /api/apps/{slug}/datatables/{name}/
HTTP Request
│
▼
AppDataTableViewSet.partial_update()
├─ _reject_if_system_table()
│
▼
DataTableSerializer.validate(attrs) same as Scenario A
│
▼
DataTableSerializer.update(instance, validated_data)
├─ SqlMapper.enrich_search_fields() if needed
├─ SqlMapper.enrich_vector_fields() idempotent safety net
└─ super().update() → DataTable.save()
├─ select_for_update() → compare old_hash vs new_hash
└─ [if schema_changed and is_materialized]
apply_migration()
└─ FrictionlessAdapter.to_sqlalchemy_table()
└─ schema_manager.apply_migration(sa_table)
└─ Alembic diff → ALTER TABLE DDL
Scenario D — Data Fetch (Normal / Vector / Hybrid)
GET /api/apps/{slug}/datatables/{name}/data/
HTTP Request + query params
│
▼
list_data() routers/data.py
├─ get_datatable_cached(app, name) Redis → DataTable
├─ resolve_allowed_actions() Cerbos authorization
├─ FilterParams.from_query_params()
│ └─ _parse_vector_value() validates float list, NaN/Inf, max dim
├─ extract_search_spec(filters, query_params)
│ ├─ strips __vector_near from filters
│ ├─ validates metric ∈ {cosine, l2, ip}, 1 ≤ topk ≤ MAX_TOPK, alpha ∈ [0,1]
│ └─ returns SearchSpec(mode, vector_field, vector_value, topk, metric, ...)
└─ provider.query_data(..., search_spec=spec)
D1 — Pure vector (spec.mode == "vector")
provider.query_data()
├─ SchemaCache.get() → SATableBuilder.build_table(schema_meta)
│ └─ [vector schema] swap vector/tsvector → string, from_descriptor,
│ write_schema(_schema_dict) → write_type(any+meta) → PgVector column
├─ _validate_search_spec() dimension mismatch → 400
├─ QueryCompiler.build(ReadTree)
│ ├─ vector_distance_expr(col, vector_value, metric)
│ │ └─ col.cosine_distance / l2_distance / max_inner_product (bound param)
│ ├─ stmt.add_columns(distance_expr.label("_vector_score"))
│ ├─ sub = stmt.subquery() distance evaluated ONCE per row (H2)
│ ├─ WHERE sub.c._vector_score <[negated if ip] threshold (H1)
│ ├─ ORDER BY sub.c._vector_score ASC
│ └─ LIMIT topk
└─ hnsw_search_session(engine, ef_search)
└─ engine.begin()
SET LOCAL hnsw.iterative_scan = relaxed_order
SET LOCAL hnsw.ef_search = N
└─ conn.execute(stmt) → rows
└─ to_similarity_score(_vector_score, metric) → _similarity_score
_vector_score vs _similarity_score (pure vector only)
Every pure-vector row carries two scores:
| Field | Meaning | Direction |
|---|---|---|
_vector_score | raw pgvector distance (debug / threshold tuning) | lower = more similar |
_similarity_score | normalised, human-facing score | higher = more similar |
_similarity_score = to_similarity_score(_vector_score, metric):
| metric | transform | range | perfect match |
|---|---|---|---|
| cosine | 1 - distance | [-1, 1] (true cosine similarity) | 1.0 |
| l2 | 1 / (1 + distance) | (0, 1] | 1.0 |
| ip | -distance (de-negate pgvector <#>) | unbounded | n/a |
ipcaveat: inner product is unbounded, so_similarity_scorefor theipmetric is the raw (de-negated) inner product — higher still means more similar, but it is not a normalised 0–1 value and1.0has no special "perfect match" meaning. Usecosineif you need a bounded similarity.
Hybrid mode does not emit
_similarity_score. Hybrid results are ranked by RRF fusion (_hybrid_score), which is a rank-based relevance score, not a distance-derived similarity — the two are not comparable, so_similarity_scoreis intentionally pure-vector only.
D2 — Hybrid vector + FTS (spec.mode == "hybrid")
HybridSearchExecutor.execute(tree, engine, provider)
├─ SATableBuilder.build_table() same as D1
│
├─ Pass 1 — Vector (alpha > 0)
│ └─ hnsw_search_session(engine)
│ SELECT pk, distance ORDER BY distance LIMIT topk×4
│
├─ Pass 2 — FTS (alpha < 1)
│ └─ engine.connect()
│ validate fts_field is TSVECTOR
│ SELECT pk, ts_rank WHERE sv @@ websearch_to_tsquery LIMIT topk×4
│
├─ RRF merge (Python)
│ score = α/(RRF_K + rank_vec) + (1−α)/(RRF_K + rank_fts)
│ sort descending, take topk
│
└─ Pass 3 — Full fetch
provider.query_data(filters={pk__in: sorted_pks})
→ adds _hybrid_score, _vector_score, _fts_score to each row
D3 — Normal relational (no vector)
provider.query_data()
└─ TreeBuilder.build(filters, sort, populate, ...)
└─ QueryCompiler.build(ReadTree)
└─ standard SA SELECT with WHERE / ORDER BY / LIMIT / LATERAL joins
Schema Cache Layer
SchemaCache.get(physical_table_name)
├─ HIT → Redis → SchemaMetadata
└─ MISS → DB (DataTable.json_schema) → build_and_cache() → Redis
SATableBuilder.build_table(schema_meta, schema_name)
├─ cache_key = "{schema_name}:{physical_table_name}:{schema_hash}"
├─ HIT → in-process dict → SA Table
└─ MISS → SqlMapper.write_schema() → SA Table → store in dict
post_save signal [on_commit]
├─ SchemaCache.invalidate(physical_table_name, schema_hash)
├─ SATableBuilder.clear_cache(schema_name, table_name) precise eviction
└─ SchemaCache.invalidate_full_package()
Quick Reference — Key Function Locations
| Function | File | Purpose |
|---|---|---|
DataTableSerializer.validate | api/serializers.py | Normalize + validate + enrich single-table schema |
DataPackageResourceSerializer.validate | api/serializers.py | Same for bulk/package endpoint |
normalize_schema_for_validation | schema/adapters/frictionless.py | vector → any, validate metadata |
enrich_vector_fields | schema/mappers/sql_mapper.py | Inject HNSW index entry into schema |
validate_vector_metadata | schema/mappers/sql_mapper.py | Validate dims/distance/index_type/storage_type |
write_type | schema/mappers/sql_mapper.py | any + x-vector-metadata → PgVector / HALFVEC |
SATableBuilder.build_table | query_engine/sa_table_builder.py | Build SA Table for read path (cached) |
extract_search_spec | api/schemas/vector_meta.py | Strip __vector_near, build SearchSpec |
vector_distance_expr | providers/flat_table/search/vector_utils.py | Metric → pgvector comparator |
hnsw_search_session | providers/flat_table/search/vector_utils.py | engine.begin() + SET LOCAL GUCs |
HybridSearchExecutor.execute | providers/flat_table/search/hybrid.py | 3-pass vector+FTS+RRF |
QueryCompiler._build_normal_select | query_engine/query_compiler.py | SA SELECT + vector subquery (H1, H2) |
DataTable.materialize | data/models.py | Create physical table via provider |
DataTable.save | data/models.py | Compute hash, set physical name, trigger materialize |