Skip to content

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/sql pulls in Thrift, which is a risk under Bun, so the REST API is called directly with fetch
  • Results always come back as EXTERNAL_LINKS chunks. INLINE is capped at 25 MiB and a single survey exceeds 500k rows. The chunk URLs are presigned, so an Authorization header must not be sent with them
  • Reads process.env directly rather than going through getEnvConfig(), because the batch scripts run without the app's required environment variables (the same reason lib/s3-client.ts read S3_URI directly)

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.

DATABRICKS_TOKEN=$(databricks auth token -p <profile> | jq -r .access_token)

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.

bun run src/batch/databricks-meta-sync.ts '{"dryRun":true}'

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.