1+ import multiprocessing
2+
13import importlib .util
4+ import json
25import os
36import urllib .parse
47from typing import Any , List , Literal , Optional , TYPE_CHECKING
58
69import pyarrow as pa
710
11+ from influxdb_client_3 .version import USER_AGENT
12+ from influxdb_client_3 .write_client ._sync import rest_client as rest
13+
814if TYPE_CHECKING :
915 import pandas as pd
1016 import polars as pl
1420from influxdb_client_3 .exceptions import InfluxDBError
1521from influxdb_client_3 .query .query_api import QueryApi as _QueryApi , QueryApiOptionsBuilder
1622from influxdb_client_3 .read_file import UploadFile
17- from influxdb_client_3 .write_client import InfluxDBClient as _InfluxDBClient , WriteOptions , Point
23+ from influxdb_client_3 .write_client import WriteOptions , Point
1824from influxdb_client_3 .write_client .client .write_api import WriteApi as _WriteApi , SYNCHRONOUS , ASYNCHRONOUS , \
1925 PointSettings , DefaultWriteOptions , WriteType
2026from influxdb_client_3 .write_client .domain .write_precision import WritePrecision
@@ -189,11 +195,16 @@ def __init__(
189195 org = None ,
190196 database = None ,
191197 token = None ,
198+ auth_scheme = None ,
199+ enable_gzip = False ,
200+ gzip_threshold = None ,
192201 write_client_options = None ,
193202 flight_client_options = None ,
194203 write_port_overwrite = None ,
195204 query_port_overwrite = None ,
196205 disable_grpc_compression = False ,
206+ point_settings = None ,
207+ debug = False ,
197208 ** kwargs ):
198209 """
199210 Initialize an InfluxDB client.
@@ -212,6 +223,14 @@ def __init__(
212223 :type flight_client_options: dict[str, any]
213224 :param disable_grpc_compression: Disable gRPC compression for Flight query responses. Default is False.
214225 :type disable_grpc_compression: bool
226+ :param point_settings The settings for Points
227+ :type point_settings: PointSettings
228+ :param debug: enable verbose logging of http requests
229+ :type debug: bool
230+ :param enable_gzip: Enable GZIP compression for write requests.
231+ :type enable_gzip: bool
232+ :param gzip_threshold: Minimum payload size (bytes) to trigger GZIP when enable_gzip is True.
233+ :type gzip_threshold: int
215234 :key auth_scheme: token authentication scheme. Set to "Bearer" for Edge.
216235 :key bool verify_ssl: Set this to false to skip verifying SSL certificate when calling API from https server.
217236 :key str ssl_ca_cert: Set this to customize the certificate file to verify the peer.
@@ -235,6 +254,10 @@ def __init__(
235254 :key bool write_no_sync: disable sync confirmation on V3 API endpoint writes.
236255 :key list[str] profilers: list of enabled Flux profilers
237256 """
257+ for key in ["host" , "token" , "database" ]:
258+ if locals ().get (key ) is None :
259+ raise Exception (f"The '{ key } ' key is required" )
260+
238261 self ._org = org if org is not None else "default"
239262 self ._database = database
240263 self ._token = token
@@ -293,14 +316,49 @@ def __init__(
293316 if write_port_overwrite is not None :
294317 port = write_port_overwrite
295318
296- self ._client = _InfluxDBClient (
297- url = f"{ scheme } ://{ hostname } :{ port } " ,
298- token = self ._token ,
319+ auth_schema = 'Token' if auth_scheme is None else auth_scheme
320+ default_header = {
321+ 'User-Agent' : USER_AGENT
322+ }
323+ if self ._token is not None :
324+ default_header ['Authorization' ] = f'{ auth_schema } { self ._token } '
325+ self .base_url = f"{ scheme } ://{ hostname } :{ port } "
326+ self .default_header = default_header
327+ self .rest_client = rest .RestClient (
328+ base_url = self .base_url ,
329+ default_header = default_header ,
330+ verify_ssl = kwargs .get ('verify_ssl' , True ),
331+ ssl_ca_cert = kwargs .get ('ssl_ca_cert' , None ),
332+ cert_file = kwargs .get ('cert_file' , None ),
333+ cert_key_file = kwargs .get ('cert_key_file' , None ),
334+ cert_key_password = kwargs .get ('cert_key_password' , None ),
335+ ssl_context = kwargs .get ('ssl_context' , None ),
336+ proxy = kwargs .get ('proxy' , None ),
337+ proxy_headers = kwargs .get ('proxy_headers' , None ),
338+ retries = kwargs .get ('retries' , False ),
339+ debug = debug ,
340+ connection_pool_maxsize = kwargs .get ('connection_pool_maxsize' , multiprocessing .cpu_count () * 5 ,)
341+ )
342+
343+ if point_settings is None :
344+ point_settings = PointSettings ()
345+
346+ # Keep WriteOptions.timeout in sync with the resolved write_timeout
347+ if isinstance (self ._write_client_options , dict ) and self ._write_client_options .get ("write_options" ) is not None :
348+ self ._write_client_options ["write_options" ].timeout = write_timeout
349+
350+ self ._write_api = _WriteApi (
351+ bucket = self ._database ,
299352 org = self ._org ,
353+ gzip_threshold = gzip_threshold ,
354+ enable_gzip = enable_gzip ,
355+ auth_scheme = auth_scheme ,
300356 timeout = write_timeout ,
301- ** kwargs )
302-
303- self ._write_api = _WriteApi (influxdb_client = self ._client , ** self ._write_client_options )
357+ default_header = default_header ,
358+ rest_client = self .rest_client ,
359+ point_settings = point_settings ,
360+ ** self ._write_client_options
361+ )
304362
305363 if query_port_overwrite is not None :
306364 port = query_port_overwrite
@@ -658,32 +716,31 @@ async def query_async(self, query: str, language: str = "sql", mode: str = "all"
658716 except ArrowException as e :
659717 raise InfluxDB3ClientQueryError (f"Error while executing query: { e } " )
660718
661- def get_server_version (self ) -> str :
719+ def get_server_version (self ) -> Optional [ str ] :
662720 """
663- Get the version of the connected InfluxDB server.
721+ Retrieves the server version by querying the designated endpoint and
722+ extracting the version information from either response headers or
723+ the response body.
664724
665- This method makes a ping request to the server and extracts the version information
666- from either the response headers or response body.
725+ This method interacts with a REST API endpoint to fetch the server's
726+ version details, which might be stored in a specific HTTP header or
727+ available in the response body as part of a JSON structure.
667728
668- :return: The version string of the InfluxDB server.
669- :rtype: str
729+ :return: The version string of the server if available, otherwise None .
730+ :rtype: Optional[ str]
670731 """
671- version = None
672- (resp_body , _ , header ) = self ._client .api_client .call_api (
673- resource_path = "/ping" ,
674- method = "GET" ,
675- response_type = object
676- )
677-
678- for key , value in header .items ():
732+ resp = self .rest_client .request (path = '/ping' , method = "GET" , headers = self .default_header )
733+ for key , value in resp .getheaders ().items ():
679734 if key .lower () == "x-influxdb-version" :
680- version = value
681- break
735+ return value
682736
683- if version is None and isinstance (resp_body , dict ):
684- version = resp_body ['version' ]
685-
686- return version
737+ try :
738+ if resp .data is not None :
739+ return json .loads (resp .data ).get ("version" )
740+ else :
741+ return None
742+ except (ValueError , TypeError ):
743+ return None
687744
688745 def flush (self ):
689746 """
@@ -700,9 +757,12 @@ def flush(self):
700757
701758 def close (self ):
702759 """Close the client and clean up resources."""
703- self ._write_api .close ()
704- self ._query_api .close ()
705- self ._client .close ()
760+ if self ._write_api is not None :
761+ self ._write_api .close ()
762+ if self ._query_api is not None :
763+ self ._query_api .close ()
764+ if self .rest_client is not None :
765+ self .rest_client .close ()
706766
707767 def __enter__ (self ):
708768 return self
0 commit comments