Maxwell
Maxwell is a Changelog Data Capture (CDC) tool that streams changes from MySQL to Apache Kafka, Kinesis, and other streaming connectors in real time. It serializes changelog events as JSON messages using a unified format.
Flink supports Maxwell JSON as both a source and a sink:
Source: Parse Maxwell JSON messages from Kafka into INSERT, UPDATE, or DELETE records in Flink SQL.
Sink: Encode Flink SQL changelog records as Maxwell JSON and write them to a data store such as Kafka.
Connectors that support the Maxwell format include the Apache Kafka connector and the Object Storage Service (OSS) connector.
Use cases
Incremental data sync: Stream database changes to another system in real time.
Log auditing: Capture and retain a record of every data modification.
Real-time materialized views: Maintain up-to-date aggregations derived from a source database table.
Temporal joins: Join a fact stream against a versioned database table using change history.
Example
The following example shows a Maxwell JSON message captured from the products table in a MySQL database. The message records an update to the row where id = 111: the weight field changed from 5.18 to 5.15.
{
"database": "test",
"table": "e",
"type": "insert",
"ts": 1477053217,
"xid": 23396,
"commit": true,
"position": "master.000006:800911",
"server_id": 23042,
"thread_id": 108,
"primary_key": [1, "2016-10-21 05:33:37.523000"],
"primary_key_columns": ["id", "c"],
"data": {
"id": 111,
"name": "scooter",
"description": "Big 2-wheel scooter",
"weight": 5.15
},
"old": {
"weight": 5.18
}
}For a description of each field, see Maxwell data format.
If this message is published to a Kafka topic named products_binlog, create the following Flink table to consume and parse it:
CREATE TABLE topic_products (
-- schema is the same as the MySQL "products" table
id BIGINT,
name STRING,
description STRING,
weight DECIMAL(10, 2)
) WITH (
'connector' = 'kafka',
'topic' = 'products_binlog',
'properties.bootstrap.servers' = 'localhost:9092',
'properties.group.id' = 'testGroup',
'format' = 'maxwell-json'
);With this table as a changelog source, run queries directly against the MySQL change stream:
-- Real-time materialized view: latest average weight per product name
SELECT name, AVG(weight) FROM topic_products GROUP BY name;
-- Sync all changes from the MySQL products table to an Elasticsearch index
INSERT INTO elasticsearch_products
SELECT * FROM topic_products;Parameters
| Parameter | Required | Default | Data type | Description |
|---|---|---|---|---|
format | Yes | (none) | STRING | The format to use. Set to maxwell-json to use Maxwell. |
maxwell-json.ignore-parse-errors | No | false | BOOLEAN | How to handle rows that fail to parse. true: skip the field or row. false: return an error and stop the deployment. |
maxwell-json.timestamp-format.standard | No | SQL | STRING | Timestamp format for input and output. SQL: parses yyyy-MM-dd HH:mm:ss.s{precision} (example: 2020-12-30 12:13:14.123). ISO-8601: parses yyyy-MM-ddTHH:mm:ss.s{precision} (example: 2020-12-30T12:13:14.123). Output uses the same format as input. |
maxwell-json.map-null-key.mode | No | FAIL | STRING | How to handle null map keys. FAIL: return an error. DROP: discard data whose key value is null in the map. LITERAL: replace null keys with the string specified by maxwell-json.map-null-key.literal. |
maxwell-json.map-null-key.literal | No | null | STRING | The string to use as a replacement for null map keys when maxwell-json.map-null-key.mode is LITERAL. |
maxwell-json.encode.decimal-as-plain-number | No | false | BOOLEAN | How to encode DECIMAL values. true: plain decimal notation (example: 0.000000027). false: scientific notation (example: 2.7E-8). |
Data type mappings
Maxwell uses JSON for serialization and deserialization. For Flink-to-JSON data type mappings, see JSON format.
Available metadata
The following metadata fields can be declared as read-only VIRTUAL columns in a DDL statement.
Format metadata is available only if the connector forwards format metadata. Only the Kafka connector supports declaring metadata fields for the value format.
| Key | Data type | Description |
|---|---|---|
database | STRING NULL | The source database. Maps to the database field in Maxwell records. |
table | STRING NULL | The source table. Maps to the table field in Maxwell records. |
primary-key-columns | ARRAY<STRING> NULL | The primary key column names. Maps to the primary_key_columns field in Maxwell records. |
ingestion-timestamp | TIMESTAMP_LTZ(3) NULL | The timestamp at which the connector processed the event. Maps to the ts field in Maxwell records. |
The following example declares all four metadata fields in a Kafka table:
CREATE TABLE KafkaTable (
origin_database STRING METADATA FROM 'value.database' VIRTUAL,
origin_table STRING METADATA FROM 'value.table' VIRTUAL,
origin_primary_key_columns ARRAY<STRING> METADATA FROM 'value.primary-key-columns' VIRTUAL,
origin_ts TIMESTAMP(3) METADATA FROM 'value.ingestion-timestamp' VIRTUAL,
user_id BIGINT,
item_id BIGINT,
behavior STRING
) WITH (
'connector' = 'kafka',
'topic' = 'user_behavior',
'properties.bootstrap.servers' = 'localhost:9092',
'properties.group.id' = 'testGroup',
'scan.startup.mode' = 'earliest-offset',
'value.format' = 'maxwell-json'
);Limitations
UPDATE encoding
Flink cannot merge an UPDATE_BEFORE message and an UPDATE_AFTER message into a single UPDATE message. When encoding Flink SQL records as Maxwell JSON:
An UPDATE_BEFORE record is encoded as a Maxwell
DELETEmessage.An UPDATE_AFTER record is encoded as a Maxwell
INSERTmessage.
Duplicate change events
In most cases, Maxwell delivers each change event exactly once (exactly-once semantics), and Flink consumes events as expected. When a fault occurs, Maxwell falls back to at-least-once semantics and may deliver duplicate events to Kafka.
To handle duplicates, set the deployment parameter table.exec.source.cdc-events-duplicate to true and define a primary key on the source. Flink then adds a stateful operator that uses the primary key to deduplicate events and produce a normalized changelog stream.