Databricks Integration
Survey answers live in Databricks Unity Catalog, not in Postgres.
Background
The application used to cache the answers.csv.gz / panels.csv.gz files that Creative Survey drops into S3 into its own Postgres tables. But a Databricks pipeline was ingesting the same files from the same bucket, so the same data was being managed twice.
The Databricks side also goes further: it joins answers with panels and splits the result per survey. Reading that directly removed the need for any cache here.
| Removed | Why |
|---|---|
s3_answers / s3_panels / s3_sync_logs |
The answer cache became unnecessary |
syncMeta / syncAnswers in batch/s3-sync.ts |
There is nothing left to sync |
lib/s3-client.ts, the S3_* variables |
S3 is no longer read |
The five /api/v1/s3/sync-* routes |
Nothing to trigger or monitor |
Tables
| Table | Purpose |
|---|---|
cs.cs_dm.dm_answers_<survey_id> |
Answers for one survey, already LEFT JOINed with panels, so no join is needed here. 611 of them exist |
cs.cs_dwh.dwh_answers / dwh_panels |
Cross-survey tables, used only to aggregate survey metadata |
cs_dm contains PII. cs_dm_pub is the same data with names, emails and phone numbers masked as ***, but the application needs exactly those fields to contact candidates, so it reads cs_dm.
The pipeline that builds them
cs_data_pipeline_s3_delivery on the Databricks side (daily at 05:00:45 JST, owned by intage).
cs_download_files_from_cs_s3 download from S3
โโโ cs_load_answers โ dl_answers
โโโ cs_load_panels โ dl_panels
โโโ cs_split_survey_tables โ cs_dm.dm_answers_<survey_id>
โโโ cs_create_non_pii_tables โ cs_dm_pub.dm_pub_answers_<survey_id>
Per-survey tables are rebuilt incrementally
cs_split_survey_tables only rebuilds surveys updated in the last three days (UPDATED_SURVEY_ID_DAYS = 3). A survey that closed more than three days ago will not have its table refreshed.
Implementation
| File | Role |
|---|---|
src/lib/databricks-client.ts |
SQL Statement Execution API client built on fetch alone |
src/repositories/databricks.ts |
Queries against the tables above |
src/batch/databricks-meta-sync.ts |
Daily refresh of databricks_surveys |
Client design
- No driver dependency.
@databricks/sqlpulls in Thrift, which is a risk under Bun, so the REST API is called directly withfetch - Results always come back as
EXTERNAL_LINKSchunks.INLINEis capped at 25 MiB and a single survey exceeds 500k rows. The chunk URLs are presigned, so anAuthorizationheader must not be sent with them - Reads
process.envdirectly rather than going throughgetEnvConfig(), because the batch scripts run without the app's required environment variables (the same reasonlib/s3-client.tsreadS3_URIdirectly)
Value types
The JSON_ARRAY format returns everything as strings.
| Column | Value |
|---|---|
is_answered / is_completed / is_mobile |
'1' / '0' |
| timestamp columns | ISO 8601 strings |
parseFlag / parseTimestamp / parseCount in repositories/databricks.ts convert them.
Handling survey_id
It is interpolated into a table name and therefore cannot be bound as a parameter. Every path goes through assertSurveyId, which only accepts digits.
Performance
| Operation | Measured |
|---|---|
| Fixed overhead per query | ~1 second |
| Full metadata aggregation (68M rows) | ~15 seconds |
| Reading one survey (531,221 rows) | ~5 seconds |
| Fetching all rows (531,221, multiple chunks) | ~165 seconds |
Because the fixed overhead dominates, batch work should fetch a lot per query. BATCH_SIZE in survey-import.ts is 1000. Measured on a 20,712 panel survey:
| Batch size | Per batch | Batches | Total |
|---|---|---|---|
| 50 | 1.53s | 415 | 10.6 min |
| 500 | 4.24s | 42 | 3.0 min |
| 1000 | 6.60s | 21 | 2.3 min |
| 2000 | 13.04s | 11 | 2.4 min |
Authentication
OAuth M2M (service principal client credentials). Tokens are cached until they expire.
DATABRICKS_HOST
DATABRICKS_WAREHOUSE_ID
DATABRICKS_CLIENT_ID
DATABRICKS_CLIENT_SECRET
DATABRICKS_CATALOG default cs
DATABRICKS_SCHEMA_DM default cs_dm
DATABRICKS_SCHEMA_DWH default cs_dwh
For local work the CLI can mint a token instead.
When DATABRICKS_TOKEN is set it is used as-is; otherwise a token is fetched through OAuth M2M.
The service principal needs Unity Catalog grants
USE_CATALOG / USE_SCHEMA / SELECT on the cs catalog. The proper home for this is membership of accs_uc_cs_cs_dm_readonly and accs_uc_cs_cs_dwh_readonly, but changing group membership is an account-level operation.
Data shape gotchas
The question identifier is the question text
The mart has no question ID. A field is addressed by its question sentence, or by the composite key question sentence|||answer item for matrix questions.
The same question text appears under several parent questions
Free-text follow-ups such as "ใใฎ็็ฑใๆใใฆใใ ใใใ" hang off several questions, so their keys collide. sub_item_sentence is empty too, leaving no way to tell them apart (only the row ID and timestamp differ).
When collapsing answers, the fields where a single panel holds more than one answer are therefore detected from the data and kept as a list (fetchMultiValueFields).
answer_type cannot tell single from multiple
The dominant type is ใใญในใ้ธๆ๏ผๅใป่ค๏ผ (42.4M rows across 8,851 questions), literally "single or multiple". Any attempt to decide from the type name is guaranteed to be wrong.
Multiple choice is stored as one row per selected option
For a multiple-choice question, each selected option gets its own row with the option text in both answer_item_sentence and value. Each option is therefore its own field.
For dropdowns and matrices, answer_item_sentence is instead a sub-question label and value is the answer. Whether answer_item_sentence === value distinguishes the two.
answer_type inventory (measured)
| answer_type | Rows | Questions |
|---|---|---|
ใใญในใ้ธๆ๏ผๅใป่ค๏ผ |
42,423,070 | 8,851 |
ใใญในใๅ
ฅๅ |
8,754,475 | 4,026 |
ใใซใใฆใณ |
8,310,220 | 664 |
ใใญในใใใใชใฏใน |
6,310,701 | 661 |
็ปๅ้ธๆ๏ผๅใป่ค๏ผ |
1,310,490 | 578 |
ในใฉใคใใใผ |
953,629 | 421 |
NPSใฟใคใ |
76,835 | 102 |
ใตใใคใใซใใญในใ |
11,867 | 9 |
ๅบๅฎๅ็ญใฟใคใ |
8,355 | 1 |
ใซใฌใณใใผ้ธๆ |
895 | 4 |
ใตใใคใใซใคใกใผใธ |
793 | 2 |
Running the batches
The ENTRYPOINT in Dockerfile.sync stops at ["bun", "run"]; the script path and its JSON argument both come from the ECS container override. One image can therefore run either batch.
# Daily survey metadata refresh
aws ecs run-task \
--cluster dev-heineken-interview-arrange-backend-sync \
--task-definition dev-heineken-interview-arrange-backend-sync \
--launch-type FARGATE \
--network-configuration "awsvpcConfiguration={subnets=[<private-subnet>],securityGroups=[<lambda-sg>],assignPublicIp=DISABLED}" \
--overrides '{"containerOverrides":[{"name":"sync","command":["src/batch/databricks-meta-sync.ts","{}"]}]}' \
--region ap-northeast-1
Locally it can be run directly.
The schedule is not wired up yet
The intent is to trigger the ECS task above after cs_data_pipeline_s3_delivery, but that hand-off does not exist. Until it does, databricks_surveys only updates when the batch is run by hand.