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

使用S3Cluster替代S3读取Parquet时ClickHouse报错求助

BigQuery通过GCS Parquet迁移至ClickHouse:S3Cluster函数报错问题

迁移流程与正常运行的脚本

BigQuery导出Parquet到GCS

使用EXPORT DATA命令导出数据,脚本正常运行:

EXPORT DATA  OPTIONS (
         uri = '{gs_uri}*{full_parquet_filename}.{gs_format.lower()}',
         format = '{gs_format}' ,
         overwrite  = true)
         AS ( SELECT * FROM `{bq_project_id}.{bq_dataset}.{bq_table}_intraday_{ts_date_str_bq_format}`
WHERE event_date = '{ts_date_str_bq_format}');

ClickHouse读取并处理嵌套列

通过S3函数读取GCS上的Parquet文件,查看列结构:

describe table s3(
    's3url',
    'access_key_id',
    'secret_access_key' 
)

其中user_properties列的类型为:

Array(Tuple(key Nullable(String), value Tuple(string_value Nullable(String), int_value Nullable(Int64), float_value Nullable(Float64), double_value Nullable(Float64), set_timestamp_micros Nullable(Int64))))

使用以下SQL将嵌套列转换为可解析的JSON字符串,该脚本通过S3函数处理数十亿行数据无问题:

select 
    arrayMap(x -> 'user_properties_'||(tupleElement(x, 1)) , user_properties ) as us_pr_key 
    ,arrayMap(x -> tupleElement(tupleElement(x, 2),1) , user_properties ) as us_pr_value_string 
    ,arrayMap(x -> tupleElement(tupleElement(x, 2),2) , user_properties ) as us_pr_value_int 
    ,arrayMap(x -> tupleElement(tupleElement(x, 2),3) , user_properties ) as us_pr_value_float
    ,arrayMap(x -> tupleElement(tupleElement(x, 2),4) , user_properties ) as us_pr_value_double

    ,arrayMap((a,b,c,d) -> 
                coalesce(toString(a),toString(b),toString(c),toString(d)) ,
                           us_pr_value_string,
                           us_pr_value_int,
                           us_pr_value_float,
                           us_pr_value_double )                                as us_pr_filled_value
    ,arrayMap((a,b) -> 
                ('{'||'"'||toString(a)||'"'||':'||'"'||toString(b)||'"'||'}') ,
                           us_pr_key,
                           us_pr_filled_value )                                as us_pr_key_value 
    ,'{'||arrayStringConcat(us_pr_key_value,', ')||'}' as us_pr_json
from 
s3(
    's3url',
    'access_key_id',
    'secret_access_key' 
)

生成的JSON示例:

{"user_properties_user_pseudo_id":"123122131241234"},{"user_properties_custom_client_id":"23123124332432"}

问题:S3Cluster函数报错

为提升插入速度,将S3替换为S3Cluster(仅添加集群名称):

..... FROM
s3Cluster(
    'clickhouse_cluster_name',
    's3url',
    'access_key_id',
    'secret_access_key' 
)

出现如下报错:

SQL Error [8] [07000]: Code: 8. DB::Exception: Received from server_name.com:9000. DB::Exception: Column 'user_properties.key' is not presented in input data.: While executing ParquetBlockInputFormat: While executing S3. (THERE_IS_NO_COLUMN) (version 23.7.3.14 (official build))
, server ClickHouseNode [uri=http://server_name.com:8123/default, options={socket_timeout=30000000,use_server_time_zone=false,use_time_zone=false}]@-1742689150

无需修改原有处理脚本的解决方案

1. 显式指定表结构

S3Cluster在分布式场景下自动解析嵌套列易出错,可通过structure参数强制指定完整表结构(从describe table s3(...)获取):

from s3Cluster(
    'clickhouse_cluster_name',
    's3url',
    'access_key_id',
    'secret_access_key',
    'Parquet',
    '-- 替换为describe table得到的完整Schema,示例:
    `user_properties` Array(Tuple(key Nullable(String), value Tuple(string_value Nullable(String), int_value Nullable(Int64), float_value Nullable(Float64), double_value Nullable(Float64), set_timestamp_micros Nullable(Int64))),
    `event_date` String,
    -- 其他列...
    '
)

此方式无需修改原有处理逻辑,仅为S3Cluster指定明确的结构解析规则。

2. 会话级启用嵌套列解析

执行查询前先设置会话参数,确保Parquet解析器正确识别嵌套结构:

SET input_format_parquet_enable_nested_columns = 1;

之后再运行原有处理脚本(替换为S3Cluster)即可,该设置不会影响其他会话。

3. 用临时视图封装原有S3查询

先基于原有S3查询创建临时视图,再通过集群查询视图实现分布式读取:

-- 创建临时视图,复用原有处理脚本
CREATE TEMPORARY VIEW temp_parquet_data AS
select 
    arrayMap(x -> 'user_properties_'||(tupleElement(x, 1)) , user_properties ) as us_pr_key 
    ,arrayMap(x -> tupleElement(tupleElement(x, 2),1) , user_properties ) as us_pr_value_string 
    ,arrayMap(x -> tupleElement(tupleElement(x, 2),2) , user_properties ) as us_pr_value_int 
    ,arrayMap(x -> tupleElement(tupleElement(x, 2),3) , user_properties ) as us_pr_value_float
    ,arrayMap(x -> tupleElement(tupleElement(x, 2),4) , user_properties ) as us_pr_value_double

    ,arrayMap((a,b,c,d) -> 
                coalesce(toString(a),toString(b),toString(c),toString(d)) ,
                           us_pr_value_string,
                           us_pr_value_int,
                           us_pr_value_float,
                           us_pr_value_double )                                as us_pr_filled_value
    ,arrayMap((a,b) -> 
                ('{'||'"'||toString(a)||'"'||':'||'"'||toString(b)||'"'||'}') ,
                           us_pr_key,
                           us_pr_filled_value )                                as us_pr_key_value 
    ,'{'||arrayStringConcat(us_pr_key_value,', ')||'}' as us_pr_json
from 
s3(
    's3url',
    'access_key_id',
    'secret_access_key' 
);

-- 分布式查询视图,实现并行读取
select * from temp_parquet_data cluster clickhouse_cluster_name;

此方式完全复用原有处理脚本,仅通过视图+集群查询实现性能提升。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 03:15:07