Updated on 2026-06-27 GMT+08:00

Enhancing the Joins of Large and Small Tables in Flink Jobs

This section applies to MRS 3.3.0 or later.

Joining Big and Small Tables

There are big tables and small tables when you join two Flink streams. Small table data is broadcasted to every join task, and large table data is rebalanced (distribute the data in a round robin fashion) to join tasks. This way, Flink SQL usability and job stability are improved.

Figure 1 Joining big and small tables

You can use Flink SQL Hints to specify the left or right table in a join as a broadcasted table and the other table as a rebalanced table. The following SQL statement examples use table A and table C as small tables:

  • Use table A as the broadcasted table.
    • Join
      SELECT /*+ BROADCAST(A) */ a2, b2 FROM A JOIN B ON a1 = b1
    • Where
      SELECT /*+ BROADCAST(A) */ a2, b2 FROM A, B WHERE a1 = b1
  • Use table A and table C as broadcasted tables.
    SELECT /*+ BROADCAST(A, C) */ a2, b2, c2 FROM A JOIN B ON a1 = b1 JOIN C ON a1 = c1
  • This feature can be used with /*+ BROADCAST(smallTable1, smallTable2) */ to be compatible with the open-source joins of two streams.
  • Switching between open-source joins and this feature is not supported because this feature broadcasts data to each join task.
  • Using a small table as the left table of a LEFTJOIN is not supported. Using a small table as the right table of a RIGHTJOIN is not supported.

However, there is a problem with this method: If the big table is updated, data may be shuffled to different partitions during the join operations because the big table uses the round-robin partitioning policy.

Specifying Field Values for Hash Partitioning (Available Only in MRS 3.6.0-LTS or Later)

In Flink SQL, you can use hash partitioning on specific field values to avoid data disorder during large table joins. If the hash field is a primary key and the large table is updated by that key, new data will be sent to the same partition as the existing data, ensuring proper data ordering during the join.

Notes and constraints

  • The specified Hash field must be the primary key.
  • The specified Hash field cannot be skewed.

How to Use

Add /*+ BROADCAST(s),HASH_BY_KEY('l'='user_id') */ to Flink SQL. BROADCAST(s) indicates that table s is a broadcast table, and HASH_BY_KEY('l'='user_id') indicates that table l uses the user_id field as the Hash field.

SELECT /*+ BROADCAST(A), HASH_BY_KEY('B'='b2') */ a2, b2 FROM A JOIN B ON a1 = b1 /* Table A is the broadcast table, and table B is hashed by column b2.*/
SELECT /*+ BROADCAST(A), HASH_BY_KEY('B'='b2, b3') */ a2, b2 FROM A JOIN B ON a1 = b1 /* Hashing can be performed based on multiple fields. However, if column b3 does not exist in table B, an error is reported.*/

Example SQL statements:

CREATE TABLE large_table (
  user_id INT,
  id INT,
  u_name STRING,
  PRIMARY KEY (user_id) NOT ENFORCED
) WITH (
 'connector' = 'hudi',
 'path' = 'hdfs://hacluster/tmp/hudi/stream_mor',
 'hoodie.datasource.write.recordkey.field' = 'user_id ',
 'read.tasks' = '1',
 'read.streaming.enabled' = 'true',
 'read.streaming.check-interval' = '5',
 'read.streaming.start-commit' = 'earliest'
);
CREATE TABLE small_table(
  id INT,
  age int
) WITH (
  'connector' = 'kafka',
  'topic' = 'small_table',
  'properties.bootstrap.servers' = 'Service IP address of the Kafka Broker instance:Kafka port',
  'scan.startup.mode' = 'latest-offset',
  'format' = 'csv'
);
SELECT
  /*+ BROADCAST(s),HASH_BY_KEY('l'='user_id') */
  l.user_id,
  s.id,
  l.u_name,
  s.age
FROM
  large_table l
  JOIN small_table s ON l.id = s.id;

Deduplicating Data When Joining Big and Small Tables

When you join two streams, there is a possibility that the join operator receives a large amount of duplicate data sent by one stream. Downstream operators need to process a large amount of duplicate data, affecting job performance.

For example, join fields (P1, A1, and A2) in table A with fields (P1, B1, B2, and B3) in table B to generate field C. A large amount of data in table B is updated but the fields remain unchanged. Assume that only the B1 and B2 fields are used in the join and only the B3 field is updated. When you update table B, B1 and B2 fields should be ignored with the deduplication function.

select  A.A1,B.B1,B.B2 from A join B on A.P1=B.P1

To deduplicate table B updates, you can use Hints to set deduplication for the left table (duplicate.left) or right table (duplicate.right).

  • Format
    • Set deduplication for the left table.
       /*+ OPTIONS('duplicate.left'='true')*/
    • Set deduplication for the right table.
       /*+ OPTIONS('duplicate.right'='true')*/
    • Set deduplication for both tables.
       /*+ OPTIONS('duplicate.left'='true','duplicate.right'='true')*/
  • The following is an example in a SQL statement:
    For example, set deduplication for both the left table user_info and the right table user_score.
    CREATE TABLE user_info (`user_id` VARCHAR, `user_name` VARCHAR) WITH (
      'connector' = 'kafka',
      'topic' = 'user_info_001',
      'properties.bootstrap.servers' = 'Service IP address of the Kafka Broker instance:Kafka port',
      'properties.group.id' = 'testGroup',
      'scan.startup.mode' = 'latest-offset',
      'value.format' = 'csv'
    );
    CREATE table print(
      `user_id` VARCHAR,
      `user_name` VARCHAR,
      `score` INT
    ) WITH ('connector' = 'print');
    CREATE TABLE user_score (user_id VARCHAR, score INT) WITH (
      'connector' = 'kafka',
      'topic' = 'user_score_001',
      'properties.bootstrap.servers' = 'Service IP address of the Kafka Broker instance:Kafka port',
      'properties.group.id' = 'testGroup',
      'scan.startup.mode' = 'latest-offset',
      'value.format' = 'csv'
    );
    INSERT INTO
      print
    SELECT
      t.user_id,
      t.user_name,
      d.score
    FROM
      user_info as t
      JOIN 
      -- Set deduplication for both left and right tables.
      user_score /*+ OPTIONS('duplicate.left'='true','duplicate.right'='true')*/ as d ON t.user_id = d.user_id;
    The IP address and port number of the Kafka Broker instance are as follows:
    • To obtain the instance IP address, log in to FusionInsight Manager, choose Cluster > Services > Kafka, click Instances, and query the instance IP address on the instance list page.
    • If Kerberos authentication is enabled for the cluster (the cluster is in security mode), the Broker port number is the value of sasl.port. The default value is 21007.
    • If Kerberos authentication is disabled for the cluster (the cluster is in normal mode), the Broker port number is the value of port. The default value is 9092. If the port number is set to 9092, set allow.everyone.if.no.acl.found to true. The procedure is as follows:

      Log in to FusionInsight Manager and choose Cluster > Services > Kafka. On the displayed page, click Configurations and then All Configurations. On the displayed page, search for allow.everyone.if.no.acl.found, set it to true, and click Save.