如何用Trino/Presto解析多行分隔日志并转换为Parquet格式?
使用Trino/Presto解析多行日志并转换为Parquet格式
核心思路
将以-----------------------为分隔符的多行日志记录,转换为结构化的行式数据,最终导出为Parquet格式。步骤如下:
- 读取日志文件并按分隔符拆分单条记录
- 将每条记录的多行键值对转换为结构化键值映射
- 将映射转换为宽表结构(变量对应列)
- 导出结果为Parquet格式
方法一:基于全文件读取(适合中小文件)
此方法将整个日志文件读取为单个字符串,再拆分处理:
-- 1. 读取日志并分割为记录块 WITH raw_data AS ( SELECT split( trim(read_text('s3://your-log-storage/path/*.log')), -- 替换为实际日志路径 '-----------------------' -- 匹配记录分隔符 ) AS record_list FROM system.one ), -- 2. 过滤空记录并拆分键值对 record_details AS ( SELECT split(trim(record), '\n') AS line_list FROM raw_data CROSS JOIN UNNEST(record_list) AS t(record) WHERE trim(record) != '' ), key_value_pairs AS ( SELECT split_part(line, '=', 1) AS var_key, split_part(line, '=', 2) AS var_value, row_number() OVER () AS record_id -- 为每条记录分配唯一ID FROM record_details CROSS JOIN UNNEST(line_list) AS t(line) WHERE trim(line) != '' ), -- 3. 聚合为宽表结构 wide_records AS ( SELECT map_agg(var_key, var_value) AS record_map FROM key_value_pairs GROUP BY record_id ) -- 4. 创建Parquet表并写入数据 CREATE TABLE your_catalog.your_target_schema.parquet_log_output WITH ( format = 'PARQUET', external_location = 's3://your-output-storage/parquet-results/', -- 替换为输出路径 partitioned_by = ARRAY[] -- 可按需添加分区字段 ) AS SELECT record_map['var1'] AS var1, record_map['var2'] AS var2, record_map['varn'] AS varn -- 列出所有需要提取的变量列 FROM wide_records;
方法二:基于逐行处理(适合大文件)
此方法通过窗口函数分组记录,避免加载整个文件到内存:
-- 1. 读取日志行并标记记录分组 WITH log_lines AS ( SELECT line, -- 统计分隔符行数量,以此分组记录 sum(CASE WHEN line LIKE '-----------------------%' THEN 1 ELSE 0 END) OVER (ORDER BY line_num) AS record_group_id FROM ( SELECT line, row_number() OVER () AS line_num FROM read_lines('s3://your-log-storage/path/*.log') -- 替换为实际路径 ) t ), -- 2. 过滤分隔符和空行,提取键值对 cleaned_key_values AS ( SELECT record_group_id, split_part(line, '=', 1) AS var_key, split_part(line, '=', 2) AS var_value FROM log_lines WHERE NOT line LIKE '-----------------------%' AND trim(line) != '' ), -- 3. 聚合为宽表 wide_records AS ( SELECT map_agg(var_key, var_value) AS record_map FROM cleaned_key_values GROUP BY record_group_id ) -- 4. 写入Parquet表 CREATE TABLE your_catalog.your_target_schema.parquet_log_output WITH ( format = 'PARQUET', external_location = 's3://your-output-storage/parquet-results/' ) AS SELECT record_map['var1'] AS var1, record_map['var2'] AS var2, record_map['varn'] AS varn FROM wide_records;
动态处理未知变量
如果日志中的变量不固定,可先查询所有唯一变量名,再动态生成SELECT语句:
- 查询所有变量:
SELECT DISTINCT var_key FROM ( SELECT split_part(line, '=', 1) AS var_key FROM read_lines('s3://your-log-storage/path/*.log') WHERE NOT line LIKE '-----------------------%' AND trim(line) != '' ) t;
- 根据查询结果,在最终SELECT子句中添加所有变量列(如
record_map['varX'] AS varX)。
内容的提问来源于stack exchange,提问作者Jason
相关产品推荐
相关产品推荐

