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
« 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
8import clickhouse_connect
9from clickhouse_connect.driver.client import Client
10from clickhouse_connect.driver.query import QueryContext
11from decouple import config
13logger = logging.getLogger(__name__)
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}
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
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 )
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"]}
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"]}
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"]
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"]
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"]
119USE_POOL = _SETTINGS["enable_connection_pool"]
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]
142 return result
144 return wrapper
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
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
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
184 return client
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")
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).")
209CH_pool: Optional[ClickHouseConnectionPool] = None
210_pool_lock = threading.Lock()
212RETRY_MAX = _SETTINGS["pool_creation_max_retries"]
213RETRY_INTERVAL = _SETTINGS["pool_creation_retries_interval"]
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)
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)
244class ClickHouseClient:
245 _client = None
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)
254 else:
255 self._client = CH_pool.get_connection()
257 self._client.execute = transform_result(self, self._client.query)
258 self._client.format = self.format
260 def __enter__(self):
261 return self._client
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
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")
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()
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