Maxwell

Updated at:

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

ParameterRequiredDefaultData typeDescription
formatYes(none)STRINGThe format to use. Set to maxwell-json to use Maxwell.
maxwell-json.ignore-parse-errorsNofalseBOOLEANHow 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.standardNoSQLSTRINGTimestamp 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.modeNoFAILSTRINGHow 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.literalNonullSTRINGThe 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-numberNofalseBOOLEANHow 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.

Important

Format metadata is available only if the connector forwards format metadata. Only the Kafka connector supports declaring metadata fields for the value format.

KeyData typeDescription
databaseSTRING NULLThe source database. Maps to the database field in Maxwell records.
tableSTRING NULLThe source table. Maps to the table field in Maxwell records.
primary-key-columnsARRAY<STRING> NULLThe primary key column names. Maps to the primary_key_columns field in Maxwell records.
ingestion-timestampTIMESTAMP_LTZ(3) NULLThe 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 DELETE message.

  • An UPDATE_AFTER record is encoded as a Maxwell INSERT message.

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.