You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何用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语句:

  1. 查询所有变量:
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;
  1. 根据查询结果,在最终SELECT子句中添加所有变量列(如record_map['varX'] AS varX)。

内容的提问来源于stack exchange,提问作者Jason

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.09 14:34:54