# kc-jdbc-src-spotify-tracks

JDBC Source Kafka connector that polls `FACTS.<env>.LABEL_SOUND_RECORDING` in Snowflake and publishes `SEARCH_TRACK_BY_ISRC` requests to the Kafka topic consumed by `lambda-store-api` (spotify-api Lambda). This triggers the Spotify track lookup pipeline for sound recordings in our catalogue.

## Pipeline

```mermaid
flowchart LR
    Snowflake_FACTS_ENV__92c8fb00[("Snowflake<br>    <br>FACTS.ENV.LABEL_SOUND_RECORDING<br><br>    Columns:<br> ISRC<br>VENDOR_ID<br>    SUBACCOUNT_ID<br> MODIFIED_AT<br>    LAST_MODIFIED_AT<br>")] -->|"JDBC poll<br>(SQL query via<br>Snowflake JDBC driver)"| kc_jdbc_src_spotify__7772d140["kc-jdbc-src-spotify-tracks (ECS Fargate)<br>    <br>Mode: timestamp<br>   <br> Watermark: COALESCE(MODIFIED_AT,<br>    LAST_MODIFIED_AT)<br>    <br>Poll: every 12h<br>    <br>SMTs: <br>ValueToKey → ExtractField$Key<br>    → ExtractField$Value<br><br>    value.converter: StringConverter"]
    kc_jdbc_src_spotify__7772d140 -->|produces to| stream_storeAPIReque_2107238e["stream.storeAPIRequest.spotify.split<br>    (MSK Kafka topic)<br>    <br>Message:<br>    {api_request: {store, type},<br>    payload: {isrc, vendor_id,<br>    subaccount_id}}"]
    stream_storeAPIReque_2107238e -->|triggers| lambda_spotify_api_A_cc79e224["lambda-spotify-api<br>    (AWS Lambda)<br>    <br>Calls Spotify Search API<br><br> GET /search?q=isrc:xyz<br>    Returns: track ID(s)<br>"]
    lambda_spotify_api_A_cc79e224 -->|produces to| event_spotifyapi_tra_43612c98["event.spotifyapi.track<br>    (MSK Kafka topic)<br>    <br>Message: track data<br>   {track.id, track.name,api_request.payload}"]
    event_spotifyapi_tra_43612c98 -->|consumed by| kc_neo_cdc_sink_spt__9fb3a83f["kc-neo-cdc-sink-spt<br>    (ECS Fargate)<br>    Runs track.cypher"]
    kc_neo_cdc_sink_spt__9fb3a83f --> Neo4j_Creates_update_5dbb8137[("Neo4j<br><br>    Creates/updates:<br>    Track:ChartmetricSpotify:PublicTrack<br>    GlobalSoundRecording<br>    <br>Sets: LabelSoundRecording.spotifyTrackId <br>(if LSR matches isrc + vendorId + subaccountId)<br>    <br>Links: GSR→Track, GSR→LSR")]
classDef style0 stroke:#730FC3,fill:#e3cff3
class Neo4j_Creates_update_5dbb8137,Snowflake_FACTS_ENV__92c8fb00,event_spotifyapi_tra_43612c98,kc_jdbc_src_spotify__7772d140,kc_neo_cdc_sink_spt__9fb3a83f,lambda_spotify_api_A_cc79e224,stream_storeAPIReque_2107238e style0
linkStyle 0,1,2,3,4,5 stroke:#788896
```

## Key configuration

| Setting | Value |
|---------|-------|
| Mode | `timestamp` |
| Watermark column | `COALESCE(MODIFIED_AT, LAST_MODIFIED_AT)::TIMESTAMP_NTZ AS WATERMARK` |
| Poll interval | 12h (`43200000ms`) |
| Topic prefix | `stream.storeAPIRequest.spotify.split` |
| Snowflake user | `{ENV}_KAFKA_CONNECT_JDBC_SOURCE_SPOTIFY_TRACKS` |
| Snowflake role | `FACTS_DB_{ENV}_SCHEMA_READ` |
| Snowflake warehouse | `{ENV}_ETL_WAREHOUSE` |
| value.converter | `StringConverter` — raw UTF-8 bytes, no JSON envelope |

## Message transforms (SMTs)

The JDBC Source connector produces a Struct record with `VALUE` and `WATERMARK` columns. Three SMTs reduce this to a raw JSON string before serialisation:

1. **ValueToKey** (`createKey`) — copies `VALUE` field to the message key
2. **ExtractField$Key** (`extractInt`) — extracts the `VALUE` field from the key struct
3. **ExtractField$Value** (`extractField`) — extracts the `VALUE` field from the value struct, leaving a bare `STRING` schema

`StringConverter` then serialises the `STRING` value as raw UTF-8 bytes. The Lambda receives `{"api_request": {"store": "SPOTIFY", "type": "SEARCH_TRACK_BY_ISRC"}, "payload": {"isrc": "...", "vendor_id": ..., "subaccount_id": ...}}`.

## Query

The connector uses a subquery wrapper so the outer `SELECT` has no `WHERE` clause. This is required because the JDBC connector appends its own `WHERE WATERMARK > ? AND WATERMARK <= ?` directly to the query string — a top-level `WHERE` would cause a SQL syntax error.

See `query.sql`.

## Secrets

| Secret | Contents |
|--------|----------|
| `{env}/kc-jdbc-src-spotify-tracks/SNOWFLAKE_PRIVATE_KEY` | `{"private_key_b64": "<base64 PEM>", "private_key_passphrase": ""}` |
| `{env}/lambda-spotify-api/SPOTIFY_CLIENT_ID` | Spotify app client ID (plaintext) |
| `{env}/lambda-spotify-api/SPOTIFY_CLIENT_SECRET` | Spotify app client secret (plaintext) |
