从HDFS迁移数据到ClickHouse频繁出现HDFS读取超时异常求助
问题现象
从HDFS向ClickHouse迁移数据时,偶尔能成功完成,但多数情况会抛出超时异常。已尝试用try-except语句实现重试、继续处理文件的逻辑,但该逻辑未生效。确认网络状态正常、文件未损坏,问题呈间歇性出现。
错误日志
[2024-03-13, 01:00:49 UTC] {logging_mixin.py:137} INFO - ClickHouse Server Exception: Code: 210. DB::Exception: 无法从HDFS读取:hdfs://eevteev:@hadoop-amber,文件路径:/user/hive/warehouse/processing.db/outbox_partitioned/dt=2024-03-10/pid=119/type=all/outbox=-1/part-00000-221487c8-dac4-4cb5-bcc9-581d496593b4.c000.txt.gz。错误:HdfsIOException: InputStreamImpl: 无法读取文件:/user/hive/warehouse/processing.db/outbox_partitioned/dt=2024-03-10/pid=119/type=all/outbox=-1/part-00000-221487c8-dac4-4cb5-bcc9-581d496593b4.c000.txt.gz,起始位置0,大小:1048576。 原因:HdfsTimeoutException: 读取8字节超时:执行ParallelParsingBlockInputFormat时出错:执行HDFSSource时出错。堆栈跟踪:
0. DB::Exception::Exception(DB::Exception::MessageMasked&&, int, bool) @ 0x000000000c7498f7 in /opt/bitnami/clickhouse/bin/clickhouse 1. DB::Exception::Exception<String&, String&, String>(int, FormatStringHelperImpl<std::type_identity<String&>::type, std::type_identity<String&>::type, std::type_identity<String>::type>, String&, String&, String&&) @ 0x0000000010d7112c in /opt/bitnami/clickhouse/bin/clickhouse 2. DB::ReadBufferFromHDFS::ReadBufferFromHDFSImpl::nextImpl() @ 0x0000000010f1c774 in /opt/bitnami/clickhouse/bin/clickhouse 3. DB::ReadBufferFromHDFS::nextImpl() @ 0x0000000010f1b6dd in /opt/bitnami/clickhouse/bin/clickhouse 4. DB::ZlibInflatingReadBuffer::nextImpl() @ 0x000000000f519939 in /opt/bitnami/clickhouse/bin/clickhouse 5. DB::segmentationEngine(DB::ReadBuffer&, DB::Memory<Allocator<false, false>>&, unsigned long, unsigned long) (.llvm.14850076424094712064) @ 0x00000000134c6b9b in /opt/bitnami/clickhouse/bin/clickhouse 6. DB::ParallelParsingInputFormat::segmentatorThreadFunction(std::shared_ptr<DB::ThreadGroup>) @ 0x00000000134fe4e0 in /opt/bitnami/clickhouse/bin/clickhouse 7. void std::__function::__policy_invoker<void ()>::__call_impl<std::__function::__default_alloc_func<ThreadFromGlobalPoolImpl<true>::ThreadFromGlobalPoolImpl<void (DB::ParallelParsingInputFormat::*)(std::shared_ptr<DB::ThreadGroup>), DB::ParallelParsingInputFormat*, std::shared_ptr<DB::ThreadGroup>>(void (DB::ParallelParsingInputFormat::*&&)(std::shared_ptr<DB::ThreadGroup>), DB::ParallelParsingInputFormat*&&, std::shared_ptr<DB::ThreadGroup>&&)::'lambda'(), void ()>>(std::__function::__policy_storage const*) @ 0x0000000013502ad6 in /opt/bitnami/clickhouse/bin/clickhouse 8. void* std::__thread_proxy[abi:v15000]<std::tuple<std::unique_ptr<std::__thread_struct, std::default_delete<std::__thread_struct>>, void ThreadPoolImpl<std::thread>::scheduleImpl(std::function<void ()>, Priority, std::optional<unsigned long>, bool)::'lambda0'()>>(void*) @ 0x000000000c832d27 in /opt/bitnami/clickhouse/bin/clickhouse 9. start_thread @ 0x0000000000007ea7 in /lib/x86_64-linux-gnu/libpthread-2.31.so 10. ? @ 0x00000000000fba2f in /lib/x86_64-linux-gnu/libc-2.31.so
现有代码
max_retries = 3 retry_count = 0 while retry_count < max_retries: for file_path in file_paths: try: result_final = client.execute(f"INSERT INTO main_gpb.yandex_gpb SELECT new_puid, sid, dt, pid FROM (SELECT CASE WHEN substring_index(line, '\\t', 1) IN (SELECT puid FROM test_gpb.gpb_cated_stream) THEN substring_index(line, '\\t', 1) ELSE TO_BASE64(substring_index(line, '\\t', 1)) END AS new_puid, arrayJoin(splitByChar(',', REGEXP_REPLACE(substring_index(line, '\\t', -1), '.*?(492708|492707|492706|492705|492704|492703|492702|492701|492700|492699|492698|492697).*', '\\\\0'))) as sid, toDate('{result_dt}') AS dt, '{pid}' AS pid FROM hdfs('hdfs://eevteev:@hadoop-amber{file_path}', 'LineAsString')) WHERE sid == '492708' or sid == '492707' or sid == '492706' or sid == '492705' or sid == '492704' or sid == '492703' or sid == '492702' or sid == '492701' or sid == '492700' or sid == '492699' or sid == '492698' or sid == '492697'") break except Exception as e: print(f"ClickHouse Server Exception: {e}") retry_count += 1 if retry_count == max_retries: print("Maximum retries reached. Exiting...") break print("Retry in 5 minutes...") time.sleep(300)
问题分析与解决方案
1. 重试逻辑失效原因
现有代码的核心问题:
- 单个文件执行失败时,
break会跳出for循环,重试时直接从下一个文件开始,不会重新尝试失败的文件。 - 重试计数器
retry_count在每个文件失败时都会递增,可能导致还没重试完当前文件就耗尽了全局重试次数。
2. 修复重试逻辑
调整代码为每个文件单独重试,确保失败的文件能被重复尝试,且不影响其他文件的处理:
import time max_retries_per_file = 3 for file_path in file_paths: retry_count = 0 success = False while retry_count < max_retries_per_file and not success: try: # 优化WHERE子句为IN,简化逻辑 result_final = client.execute(f"INSERT INTO main_gpb.yandex_gpb SELECT new_puid, sid, dt, pid FROM (SELECT CASE WHEN substring_index(line, '\\t', 1) IN (SELECT puid FROM test_gpb.gpb_cated_stream) THEN substring_index(line, '\\t', 1) ELSE TO_BASE64(substring_index(line, '\\t', 1)) END AS new_puid, arrayJoin(splitByChar(',', REGEXP_REPLACE(substring_index(line, '\\t', -1), '.*?(492708|492707|492706|492705|492704|492703|492702|492701|492700|492699|492698|492697).*', '\\\\0'))) as sid, toDate('{result_dt}') AS dt, '{pid}' AS pid FROM hdfs('hdfs://eevteev:@hadoop-amber{file_path}', 'LineAsString')) WHERE sid IN ('492708', '492707', '492706', '492705', '492704', '492703', '492702', '492701', '492700', '492699', '492698', '492697')") success = True print(f"成功处理文件: {file_path}") except Exception as e: retry_count += 1 print(f"处理文件 {file_path} 失败,重试次数 {retry_count}/{max_retries_per_file}: {e}") if retry_count < max_retries_per_file: print("等待5分钟后重试...") time.sleep(300) else: print(f"文件 {file_path} 达到最大重试次数,跳过")
3. 优化ClickHouse与HDFS的超时配置
间歇性超时大概率是连接/读取超时设置过短,可通过两种方式调整:
- 修改ClickHouse全局配置:在
config.xml中添加或更新HDFS参数<hdfs> <connect_timeout>30000</connect_timeout> <!-- 连接超时,单位毫秒 --> <read_timeout>60000</read_timeout> <!-- 读取超时,单位毫秒 --> </hdfs> - 查询语句中指定超时:直接在hdfs函数中添加参数
FROM hdfs('hdfs://eevteev:@hadoop-amber{file_path}', 'LineAsString', 'connect_timeout=30000,read_timeout=60000')
4. 其他优化建议
- 将大文件拆分为小文件,降低单次读取的数据量,减少超时概率。
- 调整ClickHouse并行解析参数
max_parallel_parsing_threads,避免线程过多导致的资源竞争。
内容的提问来源于stack exchange,提问作者howtoplay112

