how to compare (1 billion records) data between two kafka streams or Database tables

Viewed 2257

we are sending data from DB2 (table-1) via CDC to Kafka topics (topic-1). we need to do reconciliation between DB2 data and Kafka topics. we have two options -

a) bring down all kafka topic data into DB2 (as table-1-copy) and then do left outer join (between table-1 and table-1-copy) to see the non-matching records, create the delta and push it back into kafka. problem: Scalability - our data set is about a billion records and i am not sure if DB2 DBA is going to let us run such a huge join operation (that may last easily over 15-20 mins).

b) push DB2 back again into parallel kafka topic (topic-1-copy) and then do some kafka streams based solution to do left outer join between kafka topic-1 and topic-1-copy. I am still wrapping my head around kafka streams and left outer joins. I am not sure whether (using the windowing system in kafka streams) I will be able to compare the ENTIRE contents of topic-1 with topic-1-copy.

To make matters worse, the topic-1 in kafka is a compact topic, so when we push the data from DB2 back into Kafka topic-1-copy, we cannot deterministically kick off the kafka topic-compaction cycle to make sure both topic-1 and topic-1-copy are fully compacted before running any sort of compare operation on them.

c) is there any other framework option that we can consider for this ?

The ideal solution has to scale for any size data.

1 Answers

I see no reason why you couldn't do this in either Kafka Streams or KSQL. Both support table-table joins. That's assuming the format of the data is supported.

Key compaction won't affect the results as both Streams and KSQL will build the correct final state of joining the two tables. If compaction has run the amount of data that needs processing may be less, but the result will be the same.

For example, in ksqlDB you could import both topics as tables and perform a join and then filter by the topic-1 table being null to find the list of missing rows.

-- example using 0.9 ksqlDB, assuming a INT primary key:

-- create table from main topic:
CREATE TABLE_1 
   (ROWKEY INT PRIMARY KEY, <other column defs>) 
   WITH (kafka_topic='topic-1', value_format='?');

-- create table from second topic:
CREATE TABLE_2 
   (ROWKEY INT PRIMARY KEY, <other column defs>) 
   WITH (kafka_topic='topic-1-copy', value_format='?');

-- create a table containing only the missing keys:
CREATE MISSING AS
   SELECT T2.* FROM TABLE_2 T2 LEFT JOIN TABLE_1 T1
   WHERE T1.ROWKEY = null;

The benefit of this approach is that the MISSING table of missing rows would automatically update: as you extracted the missing rows from your source DB2 instance and produced them to the topic-1 then the rows in the 'MISSING' table would be deleted, i.e. you'd see tombstones being produced to the MISSING topic.

You can even extend this approach to find rows that exist in topic-1 that are no longer in the source db:

-- using the same DDL statements for TABLE_1 and TABLE_2 from above

-- perform the join:
CREATE JOINED AS
   SELECT * FROM TABLE_2 T2 FULL OUTER JOIN TABLE_1 T1;

-- detect rows in the DB that aren't in the topic:
CREATE MISSING AS
   SELECT * FROM JOINED
   WHERE T1_ROWKEY = null;

-- detect rows in the topic that aren't in the DB:
CREATE EXTRA AS
   SELECT * FROM JOINED
   WHERE T2_ROWKEY = null;

Of course, you'll need to size your cluster accordingly. The bigger your ksqlDB cluster the quicker it will process the data. It'll also need on-disk capacity to materialize the table.

The maximum amount of parallelization you can get set by the number of partitions on the topics. If you've only 1 partition, then data will be processed sequentially. If running with 100 partitions, then you can process the data using 100 CPU cores, assuming you run enough ksqlDB instances. (By default, each ksqlDB node will create 4 stream-processing threads per query, (though you can increase this if the server has more cores!)).

Related