diff --git a/CHANGELOG.md b/CHANGELOG.md index 6dcd2226..878a71e6 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -25,6 +25,10 @@ - `debug` - boolean - toggle debug level logging. - Any clients explicitly using the `Configuration`, `ApiClient`, `WriteService` or _legacy_ `InfluxDBClient` classes, will need to migrate their settings to the `InfluxDBClient3` constructor. +1. [#222](https://github.com/InfluxCommunity/influxdb3-python/pull/222): Refactor `MultiprocessingWriter` class. + - New Process will be created by using `DefaultContext.Process(target)` to prevent erratic crashes, handle cross-platform code behavior safely, and coordinate complex resource sharing. + - Users can now choose one of the start methods `fork`, `spawn` or `forkserver` when creating new Process. The default will be `spawn`. + ## 0.20.0 [2026-06-11] ### Features diff --git a/codecov.yml b/codecov.yml new file mode 100644 index 00000000..90364494 --- /dev/null +++ b/codecov.yml @@ -0,0 +1,7 @@ +coverage: + status: + project: + default: + target: auto + removed_code_behavior: fully_covered_patch + threshold: 1% \ No newline at end of file diff --git a/influxdb_client_3/write_client/client/util/multiprocessing_helper.py b/influxdb_client_3/write_client/client/util/multiprocessing_helper.py index 96ba6469..502f5c7e 100644 --- a/influxdb_client_3/write_client/client/util/multiprocessing_helper.py +++ b/influxdb_client_3/write_client/client/util/multiprocessing_helper.py @@ -7,9 +7,10 @@ import logging import multiprocessing -from influxdb_client_3 import InfluxDBClient3, write_client_options -from influxdb_client_3.write_client import WriteOptions +from influxdb_client_3 import write_client_options from influxdb_client_3.exceptions import InfluxDBError +from influxdb_client_3.write_client import WriteOptions, WriteApi +from influxdb_client_3.write_client._sync import rest_client logger = logging.getLogger('influxdb_client.client.util.multiprocessing_helper') @@ -35,9 +36,9 @@ class _PoisonPill: pass -class MultiprocessingWriter(multiprocessing.Process): +class MultiprocessingWriter: """ - The Helper class to write data into InfluxDB in independent OS process. + The Helper class to write data into InfluxDB in an independent OS process. Example: .. code-block:: python @@ -119,83 +120,80 @@ def main(): __started__ = False __disposed__ = False - def __init__(self, **kwargs) -> None: + def __init__(self, start_method='spawn', **kwargs) -> None: """ Initialize defaults. - For more information how to initialize the writer see the examples above. + For more information on how to initialize the writer, see the examples above. - :param kwargs: arguments are passed into ``__init__`` function of ``InfluxDBClient`` and ``write_api``. + :param kwargs: arguments are passed into the `` _ _init__`` function of ``InfluxDBClient`` and ``write_api``. """ - multiprocessing.Process.__init__(self) + + wco = write_client_options(write_options=kwargs.get('write_options', WriteOptions()), + success_callback=kwargs.get('success_callback', _success_callback), + error_callback=kwargs.get('error_callback', _error_callback), + retry_callback=kwargs.get('retry_callback', _retry_callback) + ) + + if kwargs.get('rest_client') is not None: + rest = kwargs.get('rest_client') + else: + token = kwargs.get('token') + default_header = {'Authorization': f'Token {token}'} + rest = rest_client.RestClient( + base_url=kwargs.get('host'), + default_header=default_header, + ) + + write_api = WriteApi( + bucket=kwargs.get('database'), + org=kwargs.get('org'), + default_header=kwargs.get('default_header'), + rest_client=rest, + **wco + ) + + self.ctx = multiprocessing.get_context(start_method) + self.process = self.ctx.Process(target=self.run, args=(write_api,)) self.kwargs = kwargs - self.client = None - self.write_api = None - self.queue_ = multiprocessing.Manager().Queue() + self.queue_ = self.ctx.JoinableQueue() def write(self, **kwargs) -> None: """ - Append time-series data into underlying queue. + Append time-series data into the underlying queue. - For more information how to pass arguments see the examples above. + For more information on how to pass arguments, see the examples above. - :param kwargs: arguments are passed into ``write`` function of ``WriteApi`` + :param kwargs: arguments are passed into the `` write `` function of ``WriteApi`` :return: None """ assert self.__disposed__ is False, 'Cannot write data: the writer is closed.' assert self.__started__ is True, 'Cannot write data: the writer is not started.' self.queue_.put(kwargs) - def run(self): + def run(self, write_api: WriteApi) -> None: """Initialize ``InfluxDBClient3`` and wait for data to write into InfluxDB.""" - # Initialize Client and Write API - wco = write_client_options(write_options=self.kwargs.get('write_options', WriteOptions()), - success_callback=self.kwargs.get('success_callback', _success_callback), - error_callback=self.kwargs.get('error_callback', _error_callback), - retry_callback=self.kwargs.get('retry_callback', _retry_callback) - ) - - # Still need to create the InfluxDBClient3 because the init logics of InfluxDBClient3 will create the WriteApi. - # it will make WriteApi class created properly. - self.client = InfluxDBClient3(write_client_options=wco, **self.kwargs) - # Close and set _query_api to None because query_api is not needed in this process. - # We only need write_api. - self.client._query_api.close() - self.client._query_api = None - - self.write_api = self.client._write_api # Infinite loop - until poison pill while True: next_record = self.queue_.get() if type(next_record) is _PoisonPill: # Poison pill means break the loop - self.terminate() + logger.info("flushing data...") + write_api.close() + logger.info("closed") self.queue_.task_done() break - self.write_api.write(**next_record) + write_api.write(**next_record) self.queue_.task_done() def start(self) -> None: - """Start independent process for writing data into InfluxDB.""" - super().start() + """Start an independent process for writing data into InfluxDB.""" + self.process.start() self.__started__ = True - def terminate(self) -> None: - """ - Cleanup resources in independent process. - - This function **cannot be used** to terminate the ``MultiprocessingWriter``. - If you want to finish your writes please call: ``__del__``. - """ - if self.write_api: - logger.info("flushing data...") - self.write_api.__del__() - self.write_api = None - if self.client: - self.client.close() - self.client = None - logger.info("closed") + def get_start_processing_method(self): + return self.ctx.get_start_method() def __enter__(self): """Enter the runtime context related to this object.""" @@ -207,11 +205,11 @@ def __exit__(self, exc_type, exc_value, traceback): self.__del__() def __del__(self): - """Dispose the client and write_api.""" + """Dispose of the client and write_api.""" if self.__started__: self.queue_.put(_PoisonPill()) self.queue_.join() - self.join() + self.process.join() self.queue_ = None self.__started__ = False self.__disposed__ = True diff --git a/pyproject.toml b/pyproject.toml index 758d2a02..d69ac4aa 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -1,3 +1,8 @@ [build-system] requires = ["setuptools>=82.0.1"] -build-backend = "setuptools.build_meta" \ No newline at end of file +build-backend = "setuptools.build_meta" + +[tool.coverage.run] +concurrency = ["multiprocessing"] +parallel = true +source = ["influxdb_client_3.write_client.client.util"] diff --git a/tests/test_influxdb_client_3_integration.py b/tests/test_influxdb_client_3_integration.py index b6225bb5..269caf2d 100644 --- a/tests/test_influxdb_client_3_integration.py +++ b/tests/test_influxdb_client_3_integration.py @@ -1,10 +1,10 @@ +import asyncio import json import logging import os import random import string import time -import asyncio import unittest import pandas as pd @@ -21,8 +21,6 @@ from influxdb_client_3.write_client.write_exceptions import ApiException from tests.util import asyncio_run, lp_to_py_object -running_on_posix = os.name == 'posix' - def random_hex(len=6): return ''.join(random.choice(string.hexdigits) for i in range(len)) @@ -343,28 +341,49 @@ def test_batch_write_closed(self): list_results = reader.to_pylist() self.assertEqual(data_size, len(list_results)) - @pytest.mark.skipif(running_on_posix, reason="Skipping this test in POSIX environments") def test_multiprocessing_helper(self): - org = 'my-org' - writer = MultiprocessingWriter( - host=self.host, - database=self.database, - token=self.token, - org=org, - write_options=WriteOptions(batch_size=1)) - writer.start() - measurement = f'test{random_hex(6)}'.lower() - for x in range(1, 10): - time.sleep(0.2) - writer.write( - bucket=self.database, - record=f"{measurement},tag=a value=\"number{x}\" {time.time_ns()}" - ) - writer.__del__() + with MultiprocessingWriter( + host=self.host, + database=self.database, + token=self.token, + org='my-org', + write_options=WriteOptions(batch_size=1)) as mp: + self.assertEqual(mp.get_start_processing_method(), 'spawn') + + measurement = f'test{random_hex(6)}'.lower() + for x in range(1, 5): + time.sleep(0.5) + mp.write( + bucket=self.database, + record=f"{measurement},tag=a value=\"number{x}\" {time.time_ns()}" + ) time.sleep(1) df = self.client.query(f'select * from {measurement}', mode="pandas") - self.assertEqual(9, len(df)) + self.assertEqual(4, len(df)) + + def test_multiprocessing_start_method_forkserver(self): + with MultiprocessingWriter( + host=self.host, + database=self.database, + token=self.token, + org='my-org', + write_options=WriteOptions(batch_size=1), + start_method='forkserver' + ) as mp: + self.assertEqual(mp.get_start_processing_method(), 'forkserver') + + def test_multiprocessing_start_method_fork(self): + + with MultiprocessingWriter( + host=self.host, + database=self.database, + token=self.token, + org='my-org', + write_options=WriteOptions(batch_size=1), + start_method='fork' + ) as mp: + self.assertEqual(mp.get_start_processing_method(), 'fork') test_cert = """-----BEGIN CERTIFICATE----- MIIDUzCCAjugAwIBAgIUZB55ULutbc9gy6xLp1BkTQU7siowDQYJKoZIhvcNAQEL