Skip to main content

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:

FieldMeaningDirection
_vector_scoreraw pgvector distance (debug / threshold tuning)lower = more similar
_similarity_scorenormalised, human-facing scorehigher = more similar

_similarity_score = to_similarity_score(_vector_score, metric):

metrictransformrangeperfect match
cosine1 - distance[-1, 1] (true cosine similarity)1.0
l21 / (1 + distance)(0, 1]1.0
ip-distance (de-negate pgvector <#>)unboundedn/a

ip caveat: inner product is unbounded, so _similarity_score for the ip metric is the raw (de-negated) inner product — higher still means more similar, but it is not a normalised 0–1 value and 1.0 has no special "perfect match" meaning. Use cosine if 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_score is 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

FunctionFilePurpose
DataTableSerializer.validateapi/serializers.pyNormalize + validate + enrich single-table schema
DataPackageResourceSerializer.validateapi/serializers.pySame for bulk/package endpoint
normalize_schema_for_validationschema/adapters/frictionless.pyvector → any, validate metadata
enrich_vector_fieldsschema/mappers/sql_mapper.pyInject HNSW index entry into schema
validate_vector_metadataschema/mappers/sql_mapper.pyValidate dims/distance/index_type/storage_type
write_typeschema/mappers/sql_mapper.pyany + x-vector-metadata → PgVector / HALFVEC
SATableBuilder.build_tablequery_engine/sa_table_builder.pyBuild SA Table for read path (cached)
extract_search_specapi/schemas/vector_meta.pyStrip __vector_near, build SearchSpec
vector_distance_exprproviders/flat_table/search/vector_utils.pyMetric → pgvector comparator
hnsw_search_sessionproviders/flat_table/search/vector_utils.pyengine.begin() + SET LOCAL GUCs
HybridSearchExecutor.executeproviders/flat_table/search/hybrid.py3-pass vector+FTS+RRF
QueryCompiler._build_normal_selectquery_engine/query_compiler.pySA SELECT + vector subquery (H1, H2)
DataTable.materializedata/models.pyCreate physical table via provider
DataTable.savedata/models.pyCompute hash, set physical name, trigger materialize