概述
FlinkSQL 提供了丰富的内置函数,覆盖了数据处理的各个方面。这些函数可以在 SQL 查询中直接使用,极大提高了数据处理的效率和灵活性。
💡 提示: 本文档基于 Flink 1.17+ 版本,某些函数在不同版本间可能有差异。
函数分类概览
| 函数类型 | 主要用途 | 常用场景 |
|---|---|---|
| 算术函数 | 数学计算 | 统计分析、指标计算 |
| 字符串函数 | 文本处理 | 数据清洗、格式转换 |
| 日期时间函数 | 时间处理 | 时间窗口、时间计算 |
| 聚合函数 | 数据汇总 | 统计报表、指标聚合 |
| 窗口函数 | 排序分析 | 排名、分组分析 |
| JSON 函数 | JSON 处理 | 半结构化数据处理 |
官方文档
详细的函数文档请参考:Apache Flink 官方文档
算术函数
用于执行数学计算的函数,支持基本运算、三角函数、对数等复杂数学运算。
基本运算
四则运算
-- 基本算术运算
SELECT
5 + 3 as addition, -- 8: 加法
10 - 4 as subtraction, -- 6: 减法
6 * 7 as multiplication, -- 42: 乘法
15 / 3 as division, -- 5: 除法
10 % 3 as modulo -- 1: 取模
FROM (VALUES (1)) AS t(x);
高级数学函数
-- 常用数学函数
SELECT
ABS(-10) as absolute_value, -- 10: 绝对值
CEIL(4.3) as ceiling, -- 5: 向上取整
FLOOR(4.7) as floor_value, -- 4: 向下取整
ROUND(4.567, 2) as rounded, -- 4.57: 四舍五入
POWER(2, 3) as power_result, -- 8: 幂运算
SQRT(16) as square_root, -- 4: 平方根
EXP(1) as exponential, -- 2.718: e的x次幂
LN(2.718) as natural_log -- 1: 自然对数
FROM (VALUES (1)) AS t(x);
三角函数
-- 三角函数示例
SELECT
SIN(PI() / 2) as sine_90, -- 1: sin(90°)
COS(0) as cosine_0, -- 1: cos(0°)
TAN(PI() / 4) as tangent_45, -- 1: tan(45°)
ASIN(1) as arcsine, -- π/2: 反正弦
ACOS(0) as arccosine, -- π/2: 反余弦
ATAN(1) as arctangent -- π/4: 反正切
FROM (VALUES (1)) AS t(x);
实际应用场景
1. 财务计算
-- 计算订单金额和税费
SELECT
order_id,
base_amount,
base_amount * 0.1 as tax_amount,
ROUND(base_amount * 1.1, 2) as total_amount,
ABS(base_amount - discount) as final_amount
FROM orders;
2. 统计分析
-- 数据标准化处理
SELECT
user_id,
score,
(score - AVG(score) OVER()) / STDDEV(score) OVER() as normalized_score,
POWER((score - 50), 2) as squared_deviation
FROM user_scores;
字符串函数
用于处理文本数据的函数,包括长度计算、子串提取、大小写转换、模式匹配等。
基本字符串操作
长度和位置
-- 字符串长度和位置函数
SELECT
CHAR_LENGTH('FlinkSQL') as char_length, -- 8: 字符长度
LENGTH('FlinkSQL') as byte_length, -- 8: 字节长度
POSITION('SQL' IN 'FlinkSQL') as position, -- 6: 子串位置
LOCATE('ink', 'FlinkSQL') as locate_pos -- 3: 查找位置
FROM (VALUES (1)) AS t(x);
子串和切割
-- 子串提取和字符串操作
SELECT
SUBSTRING('FlinkSQL', 1, 5) as substring_result, -- 'Flink': 提取子串
LEFT('FlinkSQL', 5) as left_chars, -- 'Flink': 左侧字符
RIGHT('FlinkSQL', 3) as right_chars, -- 'SQL': 右侧字符
SUBSTR('FlinkSQL', 6) as substr_from_pos -- 'SQL': 从位置开始
FROM (VALUES (1)) AS t(x);
大小写转换
-- 大小写转换函数
SELECT
UPPER('flink sql') as uppercase, -- 'FLINK SQL'
LOWER('FLINK SQL') as lowercase, -- 'flink sql'
INITCAP('flink sql streaming') as title_case -- 'Flink Sql Streaming'
FROM (VALUES (1)) AS t(x);
字符串操作和清洗
去除和填充
-- 字符串清理函数
SELECT
TRIM(' Flink ') as trimmed, -- 'Flink': 去除两端空格
LTRIM(' Flink') as left_trimmed, -- 'Flink': 去除左侧空格
RTRIM('Flink ') as right_trimmed, -- 'Flink': 去除右侧空格
LPAD('Flink', 10, '*') as left_padded, -- '****Flink': 左填充
RPAD('Flink', 10, '*') as right_padded -- 'Flink****': 右填充
FROM (VALUES (1)) AS t(x);
拼接和替换
-- 字符串拼接和替换
SELECT
CONCAT('Apache', ' ', 'Flink') as concatenated, -- 'Apache Flink'
CONCAT_WS('-', 'Flink', 'SQL', 'Stream') as joined, -- 'Flink-SQL-Stream'
REPLACE('FlinkSQL', 'SQL', 'Stream') as replaced, -- 'FlinkStream'
OVERLAY('FlinkSQL' PLACING 'Table' FROM 6) as overlaid -- 'FlinkTable'
FROM (VALUES (1)) AS t(x);
其他有用函数
-- 其他字符串函数
SELECT
REVERSE('Flink') as reversed, -- 'knilF': 反转字符串
REPEAT('Flink', 3) as repeated, -- 'FlinkFlinkFlink': 重复
CHR(65) as char_from_ascii, -- 'A': ASCII转字符
ASCII('A') as ascii_value -- 65: 字符转ASCII
FROM (VALUES (1)) AS t(x);
实际应用场景
1. 数据清洗
-- 用户输入数据清理
SELECT
user_id,
TRIM(UPPER(username)) as cleaned_username,
REPLACE(REPLACE(phone, '-', ''), ' ', '') as normalized_phone,
CONCAT(first_name, ' ', last_name) as full_name
FROM user_input;
2. 日志解析
-- 解析访问日志
SELECT
SUBSTRING(log_line, 1, 19) as timestamp_str,
SUBSTRING(log_line, POSITION(' [' IN log_line) + 2, 5) as log_level,
TRIM(SUBSTRING(log_line, POSITION('] ' IN log_line) + 2)) as message
FROM access_logs;
日期时间函数
处理日期和时间数据的函数,支持时间格式转换、日期计算、时间提取等操作。
当前时间函数
获取当前时间
-- 当前时间获取函数
SELECT
CURRENT_DATE as current_date, -- 2025-07-11: 当前日期
CURRENT_TIME as current_time, -- 14:30:45: 当前时间
CURRENT_TIMESTAMP as current_timestamp, -- 2025-07-11 14:30:45: 当前时间戳
NOW() as now_timestamp, -- 同 CURRENT_TIMESTAMP
LOCALTIMESTAMP as local_timestamp -- 本地时间戳
FROM (VALUES (1)) AS t(x);
时间格式转换
字符串与时间戳转换
-- 时间格式转换
SELECT
TO_TIMESTAMP('2025-07-11 14:30:45') as str_to_timestamp,
TO_TIMESTAMP_LTZ(1720708245, 3) as epoch_to_timestamp,
UNIX_TIMESTAMP() as current_unix_timestamp,
UNIX_TIMESTAMP('2025-07-11 14:30:45') as str_to_unix,
FROM_UNIXTIME(1720708245) as unix_to_timestamp
FROM (VALUES (1)) AS t(x);
日期格式化
-- 日期格式化函数
SELECT
DATE_FORMAT(CURRENT_TIMESTAMP, 'yyyy-MM-dd') as formatted_date,
DATE_FORMAT(CURRENT_TIMESTAMP, 'HH:mm:ss') as formatted_time,
DATE_FORMAT(CURRENT_TIMESTAMP, 'yyyy年MM月dd日') as chinese_date,
DATE_FORMAT(CURRENT_TIMESTAMP, 'E, dd MMM yyyy') as english_date
FROM (VALUES (1)) AS t(x);
日期计算
日期加减
-- 日期计算函数
SELECT
CURRENT_DATE as today,
DATE_ADD(CURRENT_DATE, INTERVAL 7 DAY) as next_week,
DATE_SUB(CURRENT_DATE, INTERVAL 1 MONTH) as last_month,
TIMESTAMPADD(HOUR, 2, CURRENT_TIMESTAMP) as two_hours_later,
TIMESTAMPDIFF(DAY, '2025-01-01', CURRENT_DATE) as days_since_new_year
FROM (VALUES (1)) AS t(x);
日期差值计算
-- 计算时间差
SELECT
order_date,
ship_date,
DATEDIFF(ship_date, order_date) as shipping_days,
TIMESTAMPDIFF(HOUR, order_time, ship_time) as shipping_hours,
TIMESTAMPDIFF(MINUTE, login_time, logout_time) as session_minutes
FROM orders;
日期部分提取
提取日期组件
-- 日期组件提取
SELECT
order_timestamp,
YEAR(order_timestamp) as order_year,
MONTH(order_timestamp) as order_month,
DAY(order_timestamp) as order_day,
HOUR(order_timestamp) as order_hour,
MINUTE(order_timestamp) as order_minute,
SECOND(order_timestamp) as order_second,
DAYOFWEEK(order_timestamp) as day_of_week, -- 1=Sunday, 7=Saturday
DAYOFYEAR(order_timestamp) as day_of_year,
WEEK(order_timestamp) as week_of_year
FROM orders;
时间窗口函数
-- 时间窗口相关函数
SELECT
event_time,
DATE_TRUNC('hour', event_time) as hour_window,
DATE_TRUNC('day', event_time) as day_window,
DATE_TRUNC('month', event_time) as month_window,
LAST_DAY(event_time) as month_end,
NEXT_DAY(event_time, 'Monday') as next_monday
FROM events;
实际应用场景
1. 时间窗口分析
-- 按小时统计订单量
SELECT
DATE_TRUNC('hour', order_time) as hour_window,
COUNT(*) as order_count,
SUM(amount) as total_amount
FROM orders
WHERE order_time >= CURRENT_DATE - INTERVAL 7 DAY
GROUP BY DATE_TRUNC('hour', order_time)
ORDER BY hour_window;
2. 业务时间计算
-- 计算用户年龄和会员时长
SELECT
user_id,
birth_date,
TIMESTAMPDIFF(YEAR, birth_date, CURRENT_DATE) as age,
register_date,
TIMESTAMPDIFF(MONTH, register_date, CURRENT_DATE) as membership_months,
CASE
WHEN TIMESTAMPDIFF(YEAR, register_date, CURRENT_DATE) >= 1
THEN 'Senior'
ELSE 'Junior'
END as membership_level
FROM users;
类型转换函数
用于在不同数据类型之间进行转换的函数。
显式类型转换
CAST 函数
-- 基本类型转换
SELECT
CAST('123' AS INT) as str_to_int, -- 123
CAST(123.45 AS INT) as float_to_int, -- 123
CAST('2025-07-11' AS DATE) as str_to_date, -- DATE '2025-07-11'
CAST(12345 AS VARCHAR) as int_to_str, -- '12345'
CAST('true' AS BOOLEAN) as str_to_bool -- TRUE
FROM (VALUES (1)) AS t(x);
TRY_CAST 安全转换
-- 安全类型转换(失败返回NULL)
SELECT
TRY_CAST('123' AS INT) as valid_conversion, -- 123
TRY_CAST('abc' AS INT) as invalid_conversion, -- NULL
TRY_CAST('2025-13-32' AS DATE) as invalid_date, -- NULL
COALESCE(TRY_CAST(input_value AS INT), 0) as safe_int
FROM input_table;
JSON 序列化
JSON 转换函数
-- JSON 序列化和反序列化
SELECT
-- 将行数据转为JSON
TO_JSON(ROW(id, name, age)) as row_to_json,
-- 解析JSON字符串
JSON_VALUE('{"name": "Alice", "age": 30}', '$.name') as json_name,
JSON_VALUE('{"name": "Alice", "age": 30}', '$.age') as json_age,
-- 检查JSON路径是否存在
JSON_EXISTS('{"name": "Alice"}', '$.age') as has_age_field
FROM users;
实际应用场景
1. 数据清洗和验证
-- 数据质量检查和清理
SELECT
raw_data,
TRY_CAST(raw_data AS INT) as cleaned_int,
CASE
WHEN TRY_CAST(raw_data AS INT) IS NULL
THEN 'INVALID'
ELSE 'VALID'
END as data_quality,
COALESCE(TRY_CAST(raw_data AS INT), -1) as default_value
FROM raw_input_table;
2. 动态数据处理
-- 处理混合类型的JSON数据
SELECT
event_id,
JSON_VALUE(properties, '$.user_id') as user_id,
CAST(JSON_VALUE(properties, '$.amount') AS DECIMAL(10,2)) as amount,
TRY_CAST(JSON_VALUE(properties, '$.timestamp') AS TIMESTAMP) as event_time
FROM event_logs
WHERE JSON_EXISTS(properties, '$.user_id');
条件函数
用于执行条件判断和逻辑处理的函数。
NULL 处理函数
NULL 判断和处理
-- NULL值处理
SELECT
user_id,
email,
email IS NULL as is_email_null,
email IS NOT NULL as has_email,
COALESCE(email, 'no-email@example.com') as email_with_default,
NULLIF(status, 'INACTIVE') as active_status,
NVL(phone, 'N/A') as phone_display
FROM users;
条件表达式
CASE WHEN 语句
-- 复杂条件判断
SELECT
user_id,
age,
CASE
WHEN age < 18 THEN 'Minor'
WHEN age BETWEEN 18 AND 65 THEN 'Adult'
WHEN age > 65 THEN 'Senior'
ELSE 'Unknown'
END as age_group,
CASE
WHEN score >= 90 THEN 'A'
WHEN score >= 80 THEN 'B'
WHEN score >= 70 THEN 'C'
WHEN score >= 60 THEN 'D'
ELSE 'F'
END as grade
FROM student_scores;
IF 函数
-- 简单条件判断
SELECT
product_id,
price,
IF(price > 100, 'Expensive', 'Affordable') as price_category,
IF(stock > 0, 'In Stock', 'Out of Stock') as availability,
IF(discount > 0, ROUND(price * (1 - discount), 2), price) as final_price
FROM products;
实际应用场景
1. 用户分级
-- 用户等级判断
SELECT
user_id,
total_orders,
total_amount,
CASE
WHEN total_amount >= 10000 THEN 'VIP'
WHEN total_amount >= 5000 THEN 'Gold'
WHEN total_amount >= 1000 THEN 'Silver'
ELSE 'Bronze'
END as membership_level,
CASE
WHEN last_login_date >= CURRENT_DATE - INTERVAL 7 DAY THEN 'Active'
WHEN last_login_date >= CURRENT_DATE - INTERVAL 30 DAY THEN 'Inactive'
ELSE 'Dormant'
END as activity_status
FROM user_summary;
2. 数据标记和分类
-- 异常检测和标记
SELECT
transaction_id,
amount,
IF(amount > 10000, 'HIGH_VALUE', 'NORMAL') as amount_flag,
CASE
WHEN transaction_time NOT BETWEEN '09:00:00' AND '18:00:00'
THEN 'OFF_HOURS'
ELSE 'BUSINESS_HOURS'
END as time_flag,
COALESCE(merchant_category, 'UNKNOWN') as category
FROM transactions;
聚合函数
用于对数据集进行统计和汇总的函数。
基本聚合函数
计数和求和
-- 基本统计函数
SELECT
COUNT(*) as total_records, -- 总记录数
COUNT(DISTINCT user_id) as unique_users, -- 唯一用户数
COUNT(email) as users_with_email, -- 非NULL邮箱数量
SUM(amount) as total_amount, -- 总金额
SUM(CASE WHEN status = 'SUCCESS' THEN 1 ELSE 0 END) as success_count
FROM orders;
平均值和极值
-- 统计分析函数
SELECT
AVG(amount) as average_amount, -- 平均值
MIN(amount) as minimum_amount, -- 最小值
MAX(amount) as maximum_amount, -- 最大值
STDDEV(amount) as standard_deviation, -- 标准差
VARIANCE(amount) as variance_value -- 方差
FROM orders;
高级聚合函数
统计分布函数
-- 高级统计函数
SELECT
department,
COUNT(*) as employee_count,
STDDEV_POP(salary) as population_stddev, -- 总体标准差
STDDEV_SAMP(salary) as sample_stddev, -- 样本标准差
VAR_POP(salary) as population_variance, -- 总体方差
VAR_SAMP(salary) as sample_variance -- 样本方差
FROM employees
GROUP BY department;
集合聚合函数
数组和映射聚合
-- 集合聚合函数
SELECT
category,
COLLECT(product_name) as product_list, -- 收集到数组
ARRAY_AGG(price) as price_array, -- 聚合到数组
MAP_AGG(product_id, product_name) as product_map, -- 聚合到映射
LISTAGG(product_name, ', ') as product_names -- 字符串聚合
FROM products
GROUP BY category;
实际应用场景
1. 销售报表
-- 每日销售统计
SELECT
DATE(order_date) as sale_date,
COUNT(*) as order_count,
COUNT(DISTINCT customer_id) as unique_customers,
SUM(amount) as daily_revenue,
AVG(amount) as average_order_value,
MIN(amount) as min_order,
MAX(amount) as max_order,
STDDEV(amount) as amount_stddev
FROM orders
WHERE order_date >= CURRENT_DATE - INTERVAL 30 DAY
GROUP BY DATE(order_date)
ORDER BY sale_date;
2. 用户行为分析
-- 用户活跃度统计
SELECT
user_id,
COUNT(*) as total_sessions,
COUNT(DISTINCT DATE(session_start)) as active_days,
SUM(session_duration) as total_duration,
AVG(session_duration) as avg_session_duration,
COLLECT(page_visited) as visited_pages,
MAX(session_start) as last_activity
FROM user_sessions
WHERE session_start >= CURRENT_DATE - INTERVAL 7 DAY
GROUP BY user_id
HAVING COUNT(*) >= 5; -- 只显示活跃用户
窗口函数
用于在查询结果集的窗口内进行分析计算的函数。
排名函数
基本排名
-- 排名函数示例
SELECT
department,
employee_name,
salary,
-- 行号(唯一递增)
ROW_NUMBER() OVER (PARTITION BY department ORDER BY salary DESC) as row_num,
-- 排名(相同值跳过)
RANK() OVER (PARTITION BY department ORDER BY salary DESC) as rank_num,
-- 密集排名(相同值不跳过)
DENSE_RANK() OVER (PARTITION BY department ORDER BY salary DESC) as dense_rank_num,
-- 百分位排名
PERCENT_RANK() OVER (PARTITION BY department ORDER BY salary DESC) as percent_rank
FROM employees;
取值函数
窗口内取值
-- 窗口取值函数
SELECT
order_date,
customer_id,
amount,
-- 窗口内第一个值
FIRST_VALUE(amount) OVER (PARTITION BY customer_id ORDER BY order_date) as first_order_amount,
-- 窗口内最后一个值
LAST_VALUE(amount) OVER (
PARTITION BY customer_id
ORDER BY order_date
ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING
) as last_order_amount,
-- 第n个值
NTH_VALUE(amount, 2) OVER (PARTITION BY customer_id ORDER BY order_date) as second_order_amount
FROM orders;
偏移函数
LAG 和 LEAD
-- 偏移函数示例
SELECT
order_date,
customer_id,
amount,
-- 前一行的值
LAG(amount, 1) OVER (PARTITION BY customer_id ORDER BY order_date) as previous_amount,
-- 后一行的值
LEAD(amount, 1) OVER (PARTITION BY customer_id ORDER BY order_date) as next_amount,
-- 计算增长率
CASE
WHEN LAG(amount, 1) OVER (PARTITION BY customer_id ORDER BY order_date) IS NOT NULL
THEN ROUND(
(amount - LAG(amount, 1) OVER (PARTITION BY customer_id ORDER BY order_date))
/ LAG(amount, 1) OVER (PARTITION BY customer_id ORDER BY order_date) * 100, 2
)
ELSE NULL
END as growth_rate
FROM orders;
聚合窗口函数
移动聚合
-- 移动窗口聚合
SELECT
order_date,
daily_sales,
-- 累计求和
SUM(daily_sales) OVER (ORDER BY order_date) as cumulative_sales,
-- 7天移动平均
AVG(daily_sales) OVER (
ORDER BY order_date
ROWS BETWEEN 6 PRECEDING AND CURRENT ROW
) as moving_avg_7day,
-- 30天移动总和
SUM(daily_sales) OVER (
ORDER BY order_date
ROWS BETWEEN 29 PRECEDING AND CURRENT ROW
) as moving_sum_30day
FROM daily_sales_summary;
实际应用场景
1. 销售排名和分析
-- 产品销售排名分析
SELECT
product_category,
product_name,
monthly_sales,
-- 类别内排名
RANK() OVER (PARTITION BY product_category ORDER BY monthly_sales DESC) as category_rank,
-- 全局排名
RANK() OVER (ORDER BY monthly_sales DESC) as global_rank,
-- 计算累计销售占比
ROUND(
SUM(monthly_sales) OVER (
PARTITION BY product_category
ORDER BY monthly_sales DESC
ROWS UNBOUNDED PRECEDING
) / SUM(monthly_sales) OVER (PARTITION BY product_category) * 100, 2
) as cumulative_percentage
FROM product_monthly_sales
WHERE sales_month = '2025-06';
2. 时间序列分析
-- 股价趋势分析
SELECT
trading_date,
stock_symbol,
closing_price,
-- 前一交易日价格
LAG(closing_price, 1) OVER (PARTITION BY stock_symbol ORDER BY trading_date) as prev_close,
-- 涨跌幅
ROUND(
(closing_price - LAG(closing_price, 1) OVER (PARTITION BY stock_symbol ORDER BY trading_date))
/ LAG(closing_price, 1) OVER (PARTITION BY stock_symbol ORDER BY trading_date) * 100, 2
) as change_percent,
-- 5日移动平均
AVG(closing_price) OVER (
PARTITION BY stock_symbol
ORDER BY trading_date
ROWS BETWEEN 4 PRECEDING AND CURRENT ROW
) as ma_5day,
-- 20日移动平均
AVG(closing_price) OVER (
PARTITION BY stock_symbol
ORDER BY trading_date
ROWS BETWEEN 19 PRECEDING AND CURRENT ROW
) as ma_20day
FROM stock_prices;
JSON 函数
用于处理 JSON 格式数据的函数。
JSON 解析函数
基本 JSON 操作
-- JSON 基本操作
SELECT
event_id,
properties,
-- 提取JSON值
JSON_VALUE(properties, '$.user_id') as user_id,
JSON_VALUE(properties, '$.amount') as amount_str,
JSON_VALUE(properties, '$.timestamp') as timestamp_str,
-- 提取并转换类型
CAST(JSON_VALUE(properties, '$.amount') AS DECIMAL(10,2)) as amount,
-- 检查JSON路径是否存在
JSON_EXISTS(properties, '$.discount') as has_discount,
-- 提取JSON对象或数组
JSON_QUERY(properties, '$.items') as items_array
FROM event_logs
WHERE JSON_EXISTS(properties, '$.user_id');
复杂 JSON 处理
-- 处理嵌套JSON数据
SELECT
order_id,
order_data,
-- 提取嵌套字段
JSON_VALUE(order_data, '$.customer.name') as customer_name,
JSON_VALUE(order_data, '$.customer.email') as customer_email,
JSON_VALUE(order_data, '$.shipping.address.city') as shipping_city,
-- 处理JSON数组
JSON_QUERY(order_data, '$.items[*].name') as item_names,
-- 计算数组长度
JSON_LENGTH(JSON_QUERY(order_data, '$.items')) as item_count
FROM orders
WHERE JSON_EXISTS(order_data, '$.customer.name');
JSON 构造函数
创建 JSON 数据
-- 构造JSON数据
SELECT
user_id,
username,
email,
-- 构造JSON对象
JSON_OBJECT(
'id', user_id,
'name', username,
'email', email,
'created_at', registration_date
) as user_json,
-- 构造JSON数组
JSON_ARRAY(user_id, username, email) as user_array
FROM users;
实际应用场景
1. 日志分析
-- 解析应用日志JSON
SELECT
DATE(log_timestamp) as log_date,
JSON_VALUE(log_data, '$.level') as log_level,
JSON_VALUE(log_data, '$.service') as service_name,
JSON_VALUE(log_data, '$.message') as log_message,
CAST(JSON_VALUE(log_data, '$.duration') AS INT) as request_duration,
COUNT(*) as log_count,
AVG(CAST(JSON_VALUE(log_data, '$.duration') AS INT)) as avg_duration
FROM application_logs
WHERE JSON_EXISTS(log_data, '$.level')
GROUP BY
DATE(log_timestamp),
JSON_VALUE(log_data, '$.level'),
JSON_VALUE(log_data, '$.service');
2. 用户行为分析
-- 分析用户行为JSON数据
SELECT
user_id,
session_date,
JSON_VALUE(behavior_data, '$.page_views') as page_views,
JSON_VALUE(behavior_data, '$.time_spent') as time_spent,
JSON_QUERY(behavior_data, '$.visited_pages[*]') as visited_pages,
JSON_LENGTH(JSON_QUERY(behavior_data, '$.visited_pages')) as unique_pages,
-- 检查特定行为
JSON_EXISTS(behavior_data, '$.actions[*] ? (@.type == "purchase")') as made_purchase
FROM user_sessions
WHERE JSON_EXISTS(behavior_data, '$.page_views');
集合函数
用于处理数组(ARRAY)和映射(MAP)等集合类型数据的函数。
数组函数
数组基本操作
-- 数组基本函数
SELECT
user_id,
favorite_categories,
-- 数组长度
CARDINALITY(favorite_categories) as category_count,
-- 检查是否包含特定元素
ARRAY_CONTAINS(favorite_categories, 'Electronics') as likes_electronics,
-- 获取数组元素
favorite_categories[1] as first_category,
-- 数组转字符串
ARRAY_JOIN(favorite_categories, ', ') as categories_str
FROM user_preferences;
数组聚合
-- 数组聚合操作
SELECT
product_category,
-- 聚合到数组
ARRAY_AGG(product_name) as product_names,
ARRAY_AGG(price ORDER BY price DESC) as prices_desc,
-- 去重聚合
ARRAY_AGG(DISTINCT brand) as unique_brands,
-- 数组大小
CARDINALITY(ARRAY_AGG(product_name)) as product_count
FROM products
GROUP BY product_category;
MAP 函数
MAP 基本操作
-- MAP 操作函数
SELECT
order_id,
product_quantities, -- MAP<STRING, INT>
-- 获取所有键
MAP_KEYS(product_quantities) as product_ids,
-- 获取所有值
MAP_VALUES(product_quantities) as quantities,
-- 获取特定键的值
product_quantities['PROD001'] as prod001_quantity,
-- 检查键是否存在
MAP_CONTAINS(product_quantities, 'PROD001') as has_prod001
FROM order_details;
MAP 聚合
-- MAP 聚合函数
SELECT
customer_id,
-- 聚合键值对到MAP
MAP_AGG(product_id, quantity) as purchase_map,
MAP_AGG(product_category, total_spent) as category_spending
FROM customer_purchases
GROUP BY customer_id;
实际应用场景
1. 用户标签分析
-- 用户标签和偏好分析
SELECT
user_id,
user_tags,
-- 标签数量
CARDINALITY(user_tags) as tag_count,
-- 检查VIP标签
ARRAY_CONTAINS(user_tags, 'VIP') as is_vip,
-- 活跃度标签
CASE
WHEN ARRAY_CONTAINS(user_tags, 'high_activity') THEN 'High'
WHEN ARRAY_CONTAINS(user_tags, 'medium_activity') THEN 'Medium'
ELSE 'Low'
END as activity_level,
-- 标签字符串
ARRAY_JOIN(user_tags, '|') as tags_string
FROM user_profiles
WHERE CARDINALITY(user_tags) > 0;
2. 商品推荐系统
-- 基于购买历史的商品推荐
WITH user_purchase_history AS (
SELECT
user_id,
ARRAY_AGG(product_category) as purchased_categories,
MAP_AGG(product_category, COUNT(*)) as category_counts
FROM purchase_history
GROUP BY user_id
)
SELECT
u.user_id,
u.purchased_categories,
-- 最常购买的类别
MAP_KEYS(u.category_counts)[1] as favorite_category,
-- 推荐新类别(排除已购买的)
ARRAY_EXCEPT(
ARRAY['Electronics', 'Clothing', 'Books', 'Sports'],
u.purchased_categories
) as recommended_categories
FROM user_purchase_history u;
正则表达式函数
用于模式匹配和文本处理的正则表达式函数。
基本正则函数
模式匹配
-- 正则表达式匹配
SELECT
user_input,
-- 检查是否匹配模式
REGEXP_LIKE(user_input, '^[A-Za-z]+$') as is_alpha_only,
REGEXP_LIKE(user_input, '^\\d{3}-\\d{3}-\\d{4}$') as is_phone_format,
REGEXP_LIKE(email, '^[\\w.-]+@[\\w.-]+\\.[A-Za-z]{2,}$') as is_valid_email,
-- 不区分大小写匹配
REGEXP_LIKE(description, 'flink', 'i') as mentions_flink
FROM user_data;
模式替换
-- 正则表达式替换
SELECT
raw_text,
-- 移除所有数字
REGEXP_REPLACE(raw_text, '\\d+', '') as text_without_numbers,
-- 规范化电话号码
REGEXP_REPLACE(phone, '[^\\d]', '') as cleaned_phone,
-- 脱敏邮箱
REGEXP_REPLACE(email, '(.{2}).*(@.*)', '$1***$2') as masked_email,
-- 替换多个空格为单个空格
REGEXP_REPLACE(description, '\\s+', ' ') as normalized_text
FROM text_data;
模式提取
-- 正则表达式提取
SELECT
log_entry,
-- 提取IP地址
REGEXP_EXTRACT(log_entry, '(\\d{1,3}\\.\\d{1,3}\\.\\d{1,3}\\.\\d{1,3})', 1) as ip_address,
-- 提取时间戳
REGEXP_EXTRACT(log_entry, '\\[(.*?)\\]', 1) as timestamp_str,
-- 提取HTTP状态码
REGEXP_EXTRACT(log_entry, '\" (\\d{3}) ', 1) as status_code,
-- 提取用户代理
REGEXP_EXTRACT(log_entry, '\"(.+?)\"$', 1) as user_agent
FROM access_logs;
实际应用场景
1. 数据清洗和验证
-- 数据质量检查和清理
SELECT
customer_id,
email,
phone,
-- 邮箱格式验证
CASE
WHEN REGEXP_LIKE(email, '^[\\w.-]+@[\\w.-]+\\.[A-Za-z]{2,}$')
THEN 'VALID'
ELSE 'INVALID'
END as email_status,
-- 电话号码清理
REGEXP_REPLACE(phone, '[^\\d]', '') as clean_phone,
-- 电话号码格式验证
CASE
WHEN REGEXP_LIKE(REGEXP_REPLACE(phone, '[^\\d]', ''), '^\\d{10,11}$')
THEN 'VALID'
ELSE 'INVALID'
END as phone_status
FROM customers;
2. 日志解析和分析
-- Web访问日志解析
SELECT
DATE(log_timestamp) as log_date,
REGEXP_EXTRACT(log_line, '(\\d{1,3}\\.\\d{1,3}\\.\\d{1,3}\\.\\d{1,3})', 1) as client_ip,
REGEXP_EXTRACT(log_line, '\"\\w+ (.+?) HTTP', 1) as requested_path,
CAST(REGEXP_EXTRACT(log_line, '\" (\\d{3}) ', 1) AS INT) as status_code,
CAST(REGEXP_EXTRACT(log_line, ' (\\d+)$', 1) AS BIGINT) as response_size,
COUNT(*) as request_count,
SUM(CAST(REGEXP_EXTRACT(log_line, ' (\\d+)$', 1) AS BIGINT)) as total_bytes
FROM web_access_logs
WHERE REGEXP_LIKE(log_line, '^\\d{1,3}\\.\\d{1,3}\\.\\d{1,3}\\.\\d{1,3}')
GROUP BY
DATE(log_timestamp),
REGEXP_EXTRACT(log_line, '(\\d{1,3}\\.\\d{1,3}\\.\\d{1,3}\\.\\d{1,3})', 1),
REGEXP_EXTRACT(log_line, '\"\\w+ (.+?) HTTP', 1),
CAST(REGEXP_EXTRACT(log_line, '\" (\\d{3}) ', 1) AS INT);
加密编码函数
用于数据加密、编码和哈希计算的函数。
哈希函数
常用哈希算法
-- 哈希函数示例
SELECT
input_data,
-- MD5 哈希
MD5(input_data) as md5_hash,
-- SHA 系列哈希
SHA1(input_data) as sha1_hash,
SHA256(input_data) as sha256_hash,
SHA512(input_data) as sha512_hash,
-- 用于密码哈希的示例
CONCAT('user_', MD5(CONCAT(user_id, 'salt_string'))) as hashed_user_id
FROM sensitive_data;
编码函数
Base64 编码
-- Base64 编码和解码
SELECT
original_text,
TO_BASE64(original_text) as base64_encoded,
FROM_BASE64(TO_BASE64(original_text)) as base64_decoded,
-- URL 编码(如果支持)
URL_ENCODE(original_text) as url_encoded
FROM text_data;
实际应用场景
1. 用户数据脱敏
-- 用户数据脱敏处理
SELECT
user_id,
-- 邮箱脱敏
CONCAT(
SUBSTRING(email, 1, 2),
'***',
SUBSTRING(email, POSITION('@' IN email))
) as masked_email,
-- 手机号脱敏
CONCAT(
SUBSTRING(phone, 1, 3),
'****',
SUBSTRING(phone, -4)
) as masked_phone,
-- 身份证号哈希
SHA256(id_number) as hashed_id,
-- 生成用户标识符
MD5(CONCAT(user_id, email, 'app_salt')) as user_token
FROM users;
2. 数据完整性校验
-- 数据完整性检查
SELECT
record_id,
data_content,
original_checksum,
SHA256(data_content) as computed_checksum,
CASE
WHEN SHA256(data_content) = original_checksum
THEN 'VALID'
ELSE 'CORRUPTED'
END as integrity_status
FROM data_records
WHERE original_checksum IS NOT NULL;
地理空间函数
用于处理地理位置和空间数据的函数。
基本空间函数
创建和操作几何对象
-- 地理空间基本操作
SELECT
location_name,
latitude,
longitude,
-- 创建点对象
ST_POINT(longitude, latitude) as point_geom,
-- 计算两点间距离(米)
ST_DISTANCE(
ST_POINT(longitude, latitude),
ST_POINT(-74.0060, 40.7128) -- 纽约坐标
) as distance_to_nyc,
-- 判断点是否在指定区域内
ST_CONTAINS(
ST_POLYGON('POLYGON((-74.1 40.6, -74.1 40.8, -73.9 40.8, -73.9 40.6, -74.1 40.6))'),
ST_POINT(longitude, latitude)
) as is_in_manhattan
FROM locations;
实际应用场景
1. 位置服务分析
-- 用户位置分析
SELECT
user_id,
ST_POINT(longitude, latitude) as user_location,
-- 查找附近的商店
COUNT(CASE
WHEN ST_DISTANCE(
ST_POINT(longitude, latitude),
ST_POINT(store_longitude, store_latitude)
) <= 1000 -- 1公里内
THEN 1
END) as nearby_stores,
-- 最近商店距离
MIN(ST_DISTANCE(
ST_POINT(longitude, latitude),
ST_POINT(store_longitude, store_latitude)
)) as nearest_store_distance
FROM user_locations u
CROSS JOIN store_locations s
GROUP BY user_id, ST_POINT(longitude, latitude);
2. 配送路径优化
-- 配送距离计算
SELECT
order_id,
pickup_location,
delivery_location,
ST_DISTANCE(pickup_location, delivery_location) as delivery_distance,
CASE
WHEN ST_DISTANCE(pickup_location, delivery_location) <= 5000
THEN 'LOCAL'
WHEN ST_DISTANCE(pickup_location, delivery_location) <= 20000
THEN 'REGIONAL'
ELSE 'LONG_DISTANCE'
END as delivery_type
FROM delivery_orders;
其他实用函数
随机数函数
随机数生成
-- 随机数函数
SELECT
-- 0-1之间的随机浮点数
RAND() as random_float,
-- 指定范围的随机整数
RAND_INTEGER(100) as random_int_0_99,
-- 随机UUID
UUID() as random_uuid,
-- 基于种子的随机数
RAND(12345) as seeded_random
FROM (VALUES (1)) AS t(x);
系统信息函数
系统和会话信息
-- 系统信息函数
SELECT
CURRENT_USER as current_user,
SESSION_USER as session_user,
VERSION() as system_version,
CURRENT_DATABASE() as current_database,
CURRENT_SCHEMA() as current_schema
FROM (VALUES (1)) AS t(x);
位运算函数
位操作
-- 位运算函数示例
SELECT
value1,
value2,
BIT_AND(value1, value2) as bit_and_result,
BIT_OR(value1, value2) as bit_or_result,
BIT_XOR(value1, value2) as bit_xor_result,
BIT_NOT(value1) as bit_not_result,
BIT_LSHIFT(value1, 2) as left_shift_2,
BIT_RSHIFT(value1, 1) as right_shift_1
FROM bit_operations_data;
实践示例
综合应用案例
1. 电商数据分析综合案例
-- 电商平台数据分析综合查询
WITH customer_analysis AS (
SELECT
c.customer_id,
c.registration_date,
c.email,
-- 客户年龄分组
CASE
WHEN TIMESTAMPDIFF(YEAR, c.birth_date, CURRENT_DATE) < 25 THEN 'Young'
WHEN TIMESTAMPDIFF(YEAR, c.birth_date, CURRENT_DATE) < 45 THEN 'Middle'
ELSE 'Senior'
END as age_group,
-- 邮箱域名提取
SUBSTRING(c.email, POSITION('@' IN c.email) + 1) as email_domain,
-- 注册时长(月)
TIMESTAMPDIFF(MONTH, c.registration_date, CURRENT_DATE) as membership_months
FROM customers c
),
order_analysis AS (
SELECT
o.customer_id,
COUNT(*) as total_orders,
SUM(o.amount) as total_spent,
AVG(o.amount) as avg_order_value,
MAX(o.order_date) as last_order_date,
-- 订单频率分析
TIMESTAMPDIFF(DAY, MIN(o.order_date), MAX(o.order_date)) / COUNT(*) as avg_days_between_orders,
-- 季节性分析
SUM(CASE WHEN QUARTER(o.order_date) = 1 THEN o.amount ELSE 0 END) as q1_spending,
SUM(CASE WHEN QUARTER(o.order_date) = 2 THEN o.amount ELSE 0 END) as q2_spending,
SUM(CASE WHEN QUARTER(o.order_date) = 3 THEN o.amount ELSE 0 END) as q3_spending,
SUM(CASE WHEN QUARTER(o.order_date) = 4 THEN o.amount ELSE 0 END) as q4_spending
FROM orders o
WHERE o.order_date >= CURRENT_DATE - INTERVAL 1 YEAR
GROUP BY o.customer_id
)
SELECT
ca.customer_id,
ca.age_group,
ca.email_domain,
ca.membership_months,
COALESCE(oa.total_orders, 0) as total_orders,
COALESCE(oa.total_spent, 0) as total_spent,
ROUND(COALESCE(oa.avg_order_value, 0), 2) as avg_order_value,
-- 客户价值分级
CASE
WHEN COALESCE(oa.total_spent, 0) >= 5000 THEN 'VIP'
WHEN COALESCE(oa.total_spent, 0) >= 1000 THEN 'Premium'
WHEN COALESCE(oa.total_spent, 0) >= 100 THEN 'Regular'
ELSE 'Basic'
END as customer_tier,
-- 活跃度评估
CASE
WHEN oa.last_order_date >= CURRENT_DATE - INTERVAL 30 DAY THEN 'Active'
WHEN oa.last_order_date >= CURRENT_DATE - INTERVAL 90 DAY THEN 'At Risk'
WHEN oa.last_order_date >= CURRENT_DATE - INTERVAL 180 DAY THEN 'Dormant'
ELSE 'Churned'
END as activity_status,
-- JSON格式的季度消费数据
JSON_OBJECT(
'Q1', COALESCE(oa.q1_spending, 0),
'Q2', COALESCE(oa.q2_spending, 0),
'Q3', COALESCE(oa.q3_spending, 0),
'Q4', COALESCE(oa.q4_spending, 0)
) as quarterly_spending
FROM customer_analysis ca
LEFT JOIN order_analysis oa ON ca.customer_id = oa.customer_id;
2. 实时日志分析案例
-- 实时日志分析和异常检测
WITH log_parsing AS (
SELECT
log_timestamp,
-- 解析IP地址
REGEXP_EXTRACT(log_line, '^(\\d{1,3}\\.\\d{1,3}\\.\\d{1,3}\\.\\d{1,3})', 1) as client_ip,
-- 解析HTTP方法和路径
REGEXP_EXTRACT(log_line, '"(\\w+) (.+?) HTTP', 1) as http_method,
REGEXP_EXTRACT(log_line, '"\\w+ (.+?) HTTP', 1) as request_path,
-- 解析状态码和响应大小
CAST(REGEXP_EXTRACT(log_line, '" (\\d{3}) ', 1) AS INT) as status_code,
CAST(REGEXP_EXTRACT(log_line, ' (\\d+)$', 1) AS BIGINT) as response_size,
-- 解析用户代理
REGEXP_EXTRACT(log_line, '"([^"]+)"$', 1) as user_agent
FROM access_logs
WHERE log_timestamp >= CURRENT_TIMESTAMP - INTERVAL 1 HOUR
),
anomaly_detection AS (
SELECT
DATE_TRUNC('minute', log_timestamp) as time_window,
client_ip,
COUNT(*) as request_count,
COUNT(DISTINCT request_path) as unique_paths,
AVG(response_size) as avg_response_size,
-- 错误率计算
ROUND(
SUM(CASE WHEN status_code >= 400 THEN 1 ELSE 0 END) * 100.0 / COUNT(*), 2
) as error_rate,
-- 检测可能的爬虫
CASE
WHEN COUNT(*) > 100 THEN 'HIGH_FREQUENCY'
WHEN COUNT(DISTINCT request_path) / COUNT(*) < 0.1 THEN 'REPETITIVE'
WHEN REGEXP_LIKE(MAX(user_agent), 'bot|crawler|spider', 'i') THEN 'BOT'
ELSE 'NORMAL'
END as traffic_type
FROM log_parsing
GROUP BY DATE_TRUNC('minute', log_timestamp), client_ip
)
SELECT
time_window,
client_ip,
request_count,
unique_paths,
ROUND(avg_response_size, 0) as avg_response_size,
error_rate,
traffic_type,
-- 异常评分
CASE
WHEN request_count > 1000 OR error_rate > 10 OR traffic_type = 'BOT' THEN 'HIGH'
WHEN request_count > 500 OR error_rate > 5 THEN 'MEDIUM'
ELSE 'LOW'
END as anomaly_score
FROM anomaly_detection
WHERE traffic_type != 'NORMAL' OR error_rate > 1
ORDER BY time_window DESC, request_count DESC;
最佳实践
1. 性能优化建议
函数使用优化
-- 推荐:在WHERE子句中避免在列上使用函数
-- 不推荐
SELECT * FROM orders WHERE YEAR(order_date) = 2025;
-- 推荐
SELECT * FROM orders
WHERE order_date >= '2025-01-01'
AND order_date < '2026-01-01';
-- 推荐:使用索引友好的查询
-- 不推荐
SELECT * FROM products WHERE UPPER(product_name) LIKE 'FLINK%';
-- 推荐
SELECT * FROM products WHERE product_name LIKE 'Flink%' OR product_name LIKE 'flink%';
聚合函数优化
-- 推荐:使用适当的聚合函数
-- 计算唯一值数量时使用COUNT(DISTINCT)而不是子查询
SELECT
category,
COUNT(DISTINCT customer_id) as unique_customers, -- 推荐
COUNT(*) as total_orders
FROM orders
GROUP BY category;
2. 数据类型选择
合适的数据类型
-- 根据数据范围选择合适的类型
SELECT
CAST(amount AS DECIMAL(10,2)) as precise_amount, -- 财务数据使用DECIMAL
CAST(user_count AS BIGINT) as user_count, -- 大数值使用BIGINT
CAST(is_active AS BOOLEAN) as is_active -- 布尔值使用BOOLEAN
FROM statistics;
3. NULL 值处理
安全的NULL处理
-- 使用COALESCE提供默认值
SELECT
user_id,
COALESCE(email, 'unknown@example.com') as email,
COALESCE(phone, 'N/A') as phone,
-- 安全的数学运算
COALESCE(score1, 0) + COALESCE(score2, 0) as total_score
FROM users;
4. 字符串处理最佳实践
字符串清理和验证
-- 标准化数据处理流程
SELECT
user_id,
-- 清理和标准化
TRIM(UPPER(COALESCE(username, ''))) as clean_username,
-- 验证邮箱格式
CASE
WHEN REGEXP_LIKE(email, '^[\\w.-]+@[\\w.-]+\\.[A-Za-z]{2,}$')
THEN email
ELSE NULL
END as validated_email,
-- 标准化电话号码
REGEXP_REPLACE(
COALESCE(phone, ''),
'[^\\d]',
''
) as clean_phone
FROM raw_user_data;
5. 时间处理最佳实践
时区和时间窗口
-- 正确处理时区和时间窗口
SELECT
-- 使用统一时区
CONVERT_TZ(event_time, 'UTC', 'America/New_York') as local_time,
-- 时间窗口边界处理
DATE_TRUNC('hour', event_time) as hour_window,
-- 业务时间判断
CASE
WHEN HOUR(event_time) BETWEEN 9 AND 17 THEN 'Business Hours'
ELSE 'Off Hours'
END as time_category
FROM events;
相关文章
系列导航:参见 Flink 系列导航
Flink SQL 核心文档
- FlinkSQL 简明教程 - FlinkSQL 基础教程和实战
- Flink Table & SQL API 实时数仓 - 数仓建设指南
Flink 核心技术
- Flink 工作流程剖析 - Flink 底层原理和架构
- Flink DataStream API - 流处理编程API
总结
FlinkSQL 内置函数提供了丰富的数据处理能力,涵盖了:
🎯 核心功能类别
- 数据计算: 算术、统计、聚合函数
- 文本处理: 字符串操作、正则表达式
- 时间处理: 日期时间函数、时间窗口
- 数据转换: 类型转换、JSON处理
- 条件逻辑: 条件判断、NULL处理
- 高级分析: 窗口函数、排名分析
💡 使用建议
- 性能优化: 避免在WHERE子句中对列使用函数
- 类型安全: 使用TRY_CAST进行安全类型转换
- NULL处理: 合理使用COALESCE和NULL判断
- 时区处理: 统一时间处理和时区转换
- 数据验证: 使用正则表达式进行格式验证
🔧 实践要点
- 根据数据类型选择合适的函数
- 注意函数的参数顺序和返回类型
- 合理使用聚合函数和窗口函数
- 重视数据清洗和验证的重要性
掌握这些内置函数,能够显著提升 FlinkSQL 的数据处理效率和查询灵活性。