更新时间:2026-06-29 GMT+08:00
分享

FlinkSQL Lookup Join支持PROCTIME()语法

使用场景

Flink Lookup Join语法中要求FOR SYSTEM_TIME AS OF包含proc_time字段,例如FOR SYSTEM_TIME AS OF o.proc_time。现在可以使用PROCTIME()代替表中的proc_time字段,直接使用FOR SYSTEM_TIME AS OF PROCTIME(),表中不需要再定义proc_time字段。子查询也适用。

约束与限制

本章节仅适用于MRS 3.6.0-LTS及之后版本。

使用方法

配置Flink作业时,Lookup Join可以直接使用PROCTIME()来代替表中proc_time字段。

SQL示例:

CREATE TABLE kafkasource(
  order_id STRING,
  price DECIMAL(32, 2),
  currency STRING
) WITH (
  'connector' = 'kafka',
  'topic' = 'test_source',
  'properties.bootstrap.servers' = 'Kafka的Broker实例业务IP:Kafka端口号',
  'properties.group.id' = 'testGroup',
  'scan.startup.mode' = 'latest-offset',
  'format' = 'csv',
  'properties.sasl.kerberos.service.name' = 'kafka',
  'properties.security.protocol' = 'SASL_PLAINTEXT',
  'properties.kerberos.domain.name' = 'hadoop.系统域名'
);
CREATE TEMPORARY TABLE customers (id STRING, name STRING, country STRING) WITH (
  'connector' = 'jdbc',
  'url' = 'jdbc:mysql://mysqlhost:ip/customerdb',
  'table-name' = 'customers',
  'username' = 'xxxx',
  'password' = 'xxxx'
);
CREATE TABLE print (
  order_id STRING,
  price DECIMAL(32, 2),
  name STRING,
  country STRING
) WITH ('connector' = 'print');
insert into
  print
SELECT
  o.order_id,
  o.price,
  c.name,
  c.country
FROM
  kafkasource as o
  JOIN customers FOR SYSTEM_TIME AS OF PROCTIME() AS c ON o.order_id = c.id;
  • Kafka Broker实例业务IP地址及端口号说明:
    • 服务的实例IP地址可通过登录Manager后,单击“集群 > 服务 > Kafka > 实例”,在实例列表页面中查询。

      登录集群Manager具体操作,请参考访问MRS集群Manager

    • 集群已启用Kerberos认证(安全模式)时Broker端口为“sasl.port”参数的值,默认为“21007”。
    • 集群未启用Kerberos认证(普通模式)时Broker端口为“port”的值,默认为“9092”。如果配置端口号为9092,则需要配置“allow.everyone.if.no.acl.found”参数为true,具体操作如下:

      登录FusionInsight Manager系统,选择“集群 > 服务 > Kafka > 配置 > 全部配置”,搜索“allow.everyone.if.no.acl.found”配置,修改参数值为true,保存配置即可。

  • 系统域名:可登录Manager,选择“系统 > 权限 > 域和互信”,查看“本端域”参数,即为当前系统域名。

相关文档