如何通过Kinesis Firehose将DynamoDB条目级变更同步至Redshift表?
Got it, let's walk through how to get your DynamoDB change data from S3 into Redshift, including handling those update records that create new entries in S3. Here's a step-by-step approach that’s worked for many folks:
Since you already have Kinesis Firehose piping DynamoDB's item-level changes to S3, you’re halfway there. The key steps now are loading that raw change data into Redshift and resolving the insert/update/delete logic properly.
1. Prepare Your Redshift Table Structure
First, you’ll need two tables: a staging table to hold raw JSON from S3, and your target table where you’ll store the cleaned, up-to-date data.
Staging Table (for raw JSON)
This table acts as a temporary holding spot for the unprocessed Firehose output:
CREATE TABLE staging_dynamodb_changes ( raw_record JSON );
Target Table (for your final data)
Match this to your DynamoDB schema, plus a timestamp to track last updates. For example, if your DynamoDB table has user_id (primary key), name, and email:
CREATE TABLE users ( user_id VARCHAR(255) PRIMARY KEY, name VARCHAR(255), email VARCHAR(255), last_updated TIMESTAMP DEFAULT CURRENT_TIMESTAMP );
2. Load Raw S3 Data into the Staging Table
Use Redshift’s COPY command to pull the JSON files from S3. Firehose typically stores files in time-partitioned prefixes (like s3://your-bucket/firehose-prefix/YYYY/MM/DD/HH/), so you can target specific batches or load everything:
COPY staging_dynamodb_changes FROM 's3://your-bucket/firehose-prefix/' IAM_ROLE 'arn:aws:iam::123456789012:role/your-redshift-access-role' FORMAT AS JSON 'auto' COMPUPDATE OFF; -- Disable compression for staging to speed up loads
If Firehose is configured to compress files (GZIP is recommended), add GZIP to the command:
COPY staging_dynamodb_changes FROM 's3://your-bucket/firehose-prefix/' IAM_ROLE 'arn:aws:iam::123456789012:role/your-redshift-access-role' FORMAT AS JSON 'auto' GZIP COMPUPDATE OFF;
3. Process Changes & Merge into the Target Table
DynamoDB’s change records include an eventName field (INSERT, MODIFY, REMOVE) and NewImage/OldImage with the item data (in DynamoDB’s typed JSON format, e.g., {"S": "string-value"}). Use Redshift’s JSON functions to parse this and handle each event type.
Handle Inserts & Updates (UPSERT)
Use Redshift’s ON CONFLICT syntax to either insert new records or update existing ones when a MODIFY event comes in:
INSERT INTO users (user_id, name, email) SELECT -- Extract values from DynamoDB's typed JSON JSON_EXTRACT_PATH_TEXT(raw_record, 'dynamodb', 'NewImage', 'user_id', 'S') AS user_id, JSON_EXTRACT_PATH_TEXT(raw_record, 'dynamodb', 'NewImage', 'name', 'S') AS name, JSON_EXTRACT_PATH_TEXT(raw_record, 'dynamodb', 'NewImage', 'email', 'S') AS email FROM staging_dynamodb_changes WHERE JSON_EXTRACT_PATH_TEXT(raw_record, 'eventName') IN ('INSERT', 'MODIFY') ON CONFLICT (user_id) DO UPDATE SET name = EXCLUDED.name, email = EXCLUDED.email, last_updated = CURRENT_TIMESTAMP;
Handle Deletes
If you need to remove records from your Redshift table when they’re deleted in DynamoDB, use the OldImage to get the primary key:
DELETE FROM users USING staging_dynamodb_changes WHERE users.user_id = JSON_EXTRACT_PATH_TEXT(staging_dynamodb_changes.raw_record, 'dynamodb', 'OldImage', 'user_id', 'S') AND JSON_EXTRACT_PATH_TEXT(staging_dynamodb_changes.raw_record, 'eventName') = 'REMOVE';
4. Automate the Workflow
To avoid manual runs, set up automation:
- Redshift Scheduled Queries: Schedule your COPY and merge SQL scripts to run on a cadence matching your Firehose batch interval (e.g., every 5 minutes).
- Lambda Trigger: Configure an S3 event trigger to run a Lambda function whenever Firehose drops new files. The function can use Redshift’s Data API to execute your loading/processing SQL.
5. Pro Tips for Smooth Operations
- Compress Firehose Output: Enable GZIP compression in your Firehose delivery stream to reduce S3 storage costs and speed up Redshift COPY operations.
- Partition S3 Data: Firehose automatically partitions files by time. Use this to load only recent batches in your COPY command (e.g.,
s3://your-bucket/firehose-prefix/2024/05/20/) to minimize data scanning. - Monitor Load Errors: Check Redshift’s
STL_LOAD_ERRORSsystem table to debug any issues with the COPY command. - Clean Up Staging Data: After processing, truncate the staging table to free up space:
TRUNCATE TABLE staging_dynamodb_changes;
内容的提问来源于stack exchange,提问作者L Xandor

