Coverage for chalicelib/utils/ch_client.py: 57%

168 statements  

« prev     ^ index     » next       coverage.py v7.15.2, created at 2026-10-10 12:56 +0000

1import logging 

2import threading 

3import time 

4from functools import wraps 

5from queue import Queue, Empty 

6from typing import Optional 

7 

8import clickhouse_connect 

9from clickhouse_connect.driver.client import Client 

10from clickhouse_connect.driver.query import QueryContext 

11from decouple import config 

12 

13logger = logging.getLogger(__name__) 

14 

15_CH_CONFIG = { 

16 "host": (config("CLICKHOUSE_HOST") if config("ch_host", default=None) is None 

17 else config("ch_host")), 

18 "user": (config("CLICKHOUSE_USER", default="default") if config("ch_user", default=None) is None 

19 else config("ch_user")), 

20 "password": (config("CLICKHOUSE_PASSWORD", default="") if config("ch_password", default=None) is None 

21 else config("ch_password")), 

22 "port": (config("CLICKHOUSE_PORT", cast=int) if config("ch_port_http", default=None) is None 

23 else config("ch_port_http", cast=int)), 

24 "client_name": config("APP_NAME", default="PY"), 

25 "database": (config("CLICKHOUSE_DATABASE", default="default") if config("ch_database", default=None) is None 

26 else config("ch_database")), 

27} 

28 

29USE_TLS = config("CLICKHOUSE_USE_TLS", cast=bool, default=False) 

30if USE_TLS: 30 ↛ 31line 30 didn't jump to line 31 because the condition on line 30 was never true

31 _CH_CONFIG["secure"] = config("CLICKHOUSE_SECURE", cast=bool, default=True) 

32 _CH_CONFIG["verify"] = ( 

33 config("CLICKHOUSE_VERIFY", cast=bool, default=True) if config("CLICKHOUSE_TLS_SKIP_VERIFY", 

34 default=None) is None 

35 else not config("CLICKHOUSE_TLS_SKIP_VERIFY", cast=bool)) 

36 _tls_cert = config("CLICKHOUSE_TLS_CERT_PATH", default=None) 

37 _tls_ca = config("CLICKHOUSE_TLS_CA_PATH", default=None) 

38 _tls_key = config("CLICKHOUSE_TLS_KEY_PATH", default=None) 

39 _tls_server_host = config("CLICKHOUSE_SERVER_HOST_NAME", default=None) 

40 if _tls_cert or _tls_ca: 

41 _CH_CONFIG["client_cert"] = _tls_cert or _tls_ca 

42 if _tls_key: 

43 _CH_CONFIG["client_cert_key"] = _tls_key 

44 if _tls_server_host: 

45 _CH_CONFIG["server_host_name"] = _tls_server_host 

46 

47CH_CONFIG = dict(_CH_CONFIG) 

48_SETTINGS = { 

49 "max_execution_time": ( 

50 config("CLICKHOUSE_MAX_EXECUTION_TIME", cast=int, default=-1) if config('ch_timeout', default=None) is None 

51 else config('ch_timeout', cast=int) 

52 ), 

53 "receive_timeout": ( 

54 config('CLICKHOUSE_RECEIVE_TIMEOUT', cast=int, default=-1) if config('ch_receive_timeout', default=None) is None 

55 else config('ch_receive_timeout', cast=int) 

56 ), 

57 "enable_compression": ( 

58 config("CLICKHOUSE_ENABLE_COMPRESSION", cast=bool, default=True) if config("CH_COMPRESSION", 

59 default=None) is None 

60 else config("CH_COMPRESSION", cast=bool) 

61 ), 

62 "compression_algorithm": config("CLICKHOUSE_COMPRESSION_ALGORITHM", default="lz4"), 

63 "query_limit": ( 

64 config("CLICKHOUSE_QUERY_LIMIT", cast=int, default=0) if config("CH_MAX_ROWS_TO_READ", default=None) is None 

65 else config("CH_MAX_ROWS_TO_READ", cast=int) 

66 ), 

67 "query_retries": ( 

68 config("CLICKHOUSE_QUERY_RETRIES", cast=int, default=0) if config("CH_QUERY_RETRIES", default=None) is None 

69 else config("CH_QUERY_RETRIES", cast=int) 

70 ), 

71 "enable_connection_pool": ( 

72 config('CLICKHOUSE_USE_CONNECTION_POOL', cast=bool, default=True) if config('CH_POOL', default=None) is None 

73 else config('CH_POOL', cast=bool) 

74 ), 

75 "pool_min_connections": ( 

76 config("CLICKHOUSE_POOL_MIN_CONNECTIONS", cast=int, default=4) if config("CH_MINCONN", default=None) is None 

77 else config("CH_MINCONN", cast=int) 

78 ), 

79 "pool_max_connections": ( 

80 config("CLICKHOUSE_POOL_MAX_CONNECTIONS", cast=int, default=8) if config("CH_MAXCONN", default=None) is None 

81 else config("CH_MAXCONN", cast=int) 

82 ), 

83 "pool_get_timeout": ( 

84 config("CLICKHOUSE_GET_CNX_POOL_TIMEOUT_S", cast=int, default=10) if config("CH_WAIT_FOR_CNX_POOL_S", 

85 default=None) is None 

86 else config("CH_WAIT_FOR_CNX_POOL_S", cast=int) 

87 ), 

88 "pool_creation_max_retries": ( 

89 config("CLICKHOUSE_POOL_CREATE_MAX_RETRIES", cast=int, default=10) if config("CH_RETRY_MAX", 

90 default=None) is None 

91 else config("CH_RETRY_MAX", cast=int) 

92 ), 

93 "pool_creation_retries_interval": ( 

94 config("CLICKHOUSE_POOL_CREATION_RETRIES_INTERVAL_S", cast=int, default=2) if config("CH_RETRY_INTERVAL", 

95 default=None) is None 

96 else config("CH_RETRY_INTERVAL", cast=int, default=2) 

97 ) 

98 

99} 

100settings = {} 

101if _SETTINGS["max_execution_time"] > 0: 101 ↛ 105line 101 didn't jump to line 105 because the condition on line 101 was always true

102 logging.info(f"CH-max_execution_time set to {_SETTINGS["max_execution_time"]}s") 

103 settings = {**settings, "max_execution_time": _SETTINGS["max_execution_time"]} 

104 

105if _SETTINGS["receive_timeout"] > 0: 105 ↛ 109line 105 didn't jump to line 109 because the condition on line 105 was always true

106 logging.info(f"CH-receive_timeout set to {_SETTINGS["receive_timeout"]}s") 

107 settings = {**settings, "receive_timeout": _SETTINGS["receive_timeout"]} 

108 

109extra_args = {} 

110if _SETTINGS["enable_compression"]: 110 ↛ 113line 110 didn't jump to line 113 because the condition on line 110 was always true

111 extra_args["compression"] = _SETTINGS["compression_algorithm"] 

112 

113if _SETTINGS["query_limit"] > 0: 113 ↛ 116line 113 didn't jump to line 116 because the condition on line 113 was always true

114 extra_args["query_limit"] = _SETTINGS["query_limit"] 

115 

116if _SETTINGS["query_retries"] > 0: 116 ↛ 119line 116 didn't jump to line 119 because the condition on line 116 was always true

117 extra_args["query_retries"] = _SETTINGS["query_retries"] 

118 

119USE_POOL = _SETTINGS["enable_connection_pool"] 

120 

121 

122def transform_result(self, original_function): 

123 @wraps(original_function) 

124 def wrapper(*args, **kwargs): 

125 if kwargs.get("parameters"): 

126 if config("LOCAL_DEV", cast=bool, default=False): 126 ↛ 127line 126 didn't jump to line 127 because the condition on line 126 was never true

127 logger.debug(self.format(query=kwargs.get("query", ""), parameters=kwargs.get("parameters"))) 

128 else: 

129 logger.debug( 

130 str.encode(self.format(query=kwargs.get("query", ""), parameters=kwargs.get("parameters")))) 

131 elif len(args) > 0: 

132 if config("LOCAL_DEV", cast=bool, default=False): 132 ↛ 133line 132 didn't jump to line 133 because the condition on line 132 was never true

133 logger.debug(args[0]) 

134 else: 

135 logger.debug(str.encode(args[0])) 

136 result = original_function(*args, **kwargs) 

137 if isinstance(result, clickhouse_connect.driver.query.QueryResult): 137 ↛ 142line 137 didn't jump to line 142 because the condition on line 137 was always true

138 column_names = result.column_names 

139 result = result.result_rows 

140 result = [dict(zip(column_names, row)) for row in result] 

141 

142 return result 

143 

144 return wrapper 

145 

146 

147class ClickHouseConnectionPool: 

148 def __init__(self, min_size, max_size): 

149 self.min_size = min_size 

150 self.max_size = max_size 

151 self.pool: Queue[Client] = Queue() 

152 self.lock = threading.Lock() 

153 self.total_connections = 0 

154 

155 # Initialize the pool with min_size connections 

156 for _ in range(self.min_size): 

157 client = clickhouse_connect.get_client(**CH_CONFIG, 

158 settings=settings, 

159 **extra_args) 

160 self.pool.put(client) 

161 self.total_connections += 1 

162 

163 def get_connection(self): 

164 try: 

165 # Try to get a connection without blocking 

166 client = self.pool.get_nowait() 

167 except Empty as e: 

168 with self.lock: 

169 if self.total_connections < self.max_size: 

170 # If no connexion in the pool is found, create a new one if max not reached 

171 client = clickhouse_connect.get_client(**CH_CONFIG, 

172 settings=settings, 

173 **extra_args) 

174 self.total_connections += 1 

175 else: 

176 # If max_size reached, wait until a connection is available 

177 logger.info("Total connections exceeded, waiting for a new connection in the pool") 

178 try: 

179 client = self.pool.get(timeout=_SETTINGS["pool_get_timeout"]) 

180 except Empty as exc: 

181 logger.error("Pool wait timeout exceeded, no connections left") 

182 raise exc 

183 

184 return client 

185 

186 def release_connection(self, client): 

187 # If the queue is full (should not happen with consistent get/release), 

188 # close the extra connection to avoid leaks. 

189 try: 

190 self.pool.put_nowait(client) 

191 except Exception: 

192 try: 

193 client.close() 

194 except Exception: 

195 logger.exception("Error while closing extra ClickHouse client") 

196 

197 def close_all(self): 

198 with self.lock: 

199 while not self.pool.empty(): 

200 client = self.pool.get() 

201 try: 

202 client.close() 

203 except Exception: 

204 logger.exception("Error while closing ClickHouse client") 

205 self.total_connections = 0 

206 logger.info("Closed all ClickHouse connections (pool).") 

207 

208 

209CH_pool: Optional[ClickHouseConnectionPool] = None 

210_pool_lock = threading.Lock() 

211 

212RETRY_MAX = _SETTINGS["pool_creation_max_retries"] 

213RETRY_INTERVAL = _SETTINGS["pool_creation_retries_interval"] 

214 

215 

216def make_pool(): 

217 if not USE_POOL: 217 ↛ 218line 217 didn't jump to line 218 because the condition on line 217 was never true

218 return 

219 global CH_pool 

220 with _pool_lock: 

221 if CH_pool is not None: 221 ↛ 222line 221 didn't jump to line 222 because the condition on line 221 was never true

222 try: 

223 CH_pool.close_all() 

224 except Exception as error: 

225 logger.error("Error while closing all connexions to CH", exc_info=error) 

226 

227 attempts = 0 

228 while True: 

229 attempts += 1 

230 try: 

231 CH_pool = ClickHouseConnectionPool(min_size=_SETTINGS["pool_min_connections"], 

232 max_size=_SETTINGS["pool_max_connections"]) 

233 logger.info("ClickHouse connection pool created successfully") 

234 return 

235 except Exception as error: 

236 logger.error("Error while creating ClickHouse pool", exc_info=error) 

237 if attempts >= RETRY_MAX: 

238 logger.error("Failed to create ClickHouse pool after %d attempts", attempts) 

239 raise 

240 logger.info(f"waiting for {RETRY_INTERVAL}s before retry n°{attempts}") 

241 time.sleep(RETRY_INTERVAL) 

242 

243 

244class ClickHouseClient: 

245 _client = None 

246 

247 def __init__(self): 

248 if self._client is None: 248 ↛ exitline 248 didn't return from function '__init__' because the condition on line 248 was always true

249 if not USE_POOL: 249 ↛ 250line 249 didn't jump to line 250 because the condition on line 249 was never true

250 self._client = clickhouse_connect.get_client(**CH_CONFIG, 

251 settings=settings, 

252 **extra_args) 

253 

254 else: 

255 self._client = CH_pool.get_connection() 

256 

257 self._client.execute = transform_result(self, self._client.query) 

258 self._client.format = self.format 

259 

260 def __enter__(self): 

261 return self._client 

262 

263 @staticmethod 

264 def format(query, parameters=None): 

265 if parameters: 265 ↛ 268line 265 didn't jump to line 268 because the condition on line 265 was always true

266 ctx = QueryContext(query=query, parameters=parameters) 

267 return ctx.final_query 

268 return query 

269 

270 def __exit__(self, *args): 

271 if USE_POOL: 271 ↛ 274line 271 didn't jump to line 274 because the condition on line 271 was always true

272 CH_pool.release_connection(self._client) 

273 else: 

274 try: 

275 self._client.close() 

276 except Exception: 

277 logger.exception("Error while closing ClickHouse client") 

278 

279 

280async def init() -> None: 

281 logger.info(f">use CH_POOL:{USE_POOL}") 

282 if USE_POOL: 282 ↛ exitline 282 didn't return from function 'init' because the condition on line 282 was always true

283 make_pool() 

284 

285 

286async def terminate() -> None: 

287 global CH_pool 

288 if CH_pool is not None: 

289 try: 

290 CH_pool.close_all() 

291 logger.info("Closed all connexions to CH") 

292 except Exception as error: 

293 logger.error("Error while closing all connexions to CH", exc_info=error) 

294 finally: 

295 CH_pool = None