使用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
相关产品推荐
相关产品推荐

