airbyte.logs

PyAirbyte Logging features and related configuration.

By default, PyAirbyte main logs are written to a file in the AIRBYTE_LOGGING_ROOT directory, which defaults to a system-created temporary directory. PyAirbyte also maintains connector-specific log files within the same directory, under a subfolder with the name of the connector.

PyAirbyte supports structured JSON logging, which is disabled by default. To enable structured logging in JSON, set AIRBYTE_STRUCTURED_LOGGING to True.

  1# Copyright (c) 2024 Airbyte, Inc., all rights reserved.
  2"""PyAirbyte Logging features and related configuration.
  3
  4By default, PyAirbyte main logs are written to a file in the `AIRBYTE_LOGGING_ROOT` directory, which
  5defaults to a system-created temporary directory. PyAirbyte also maintains connector-specific log
  6files within the same directory, under a subfolder with the name of the connector.
  7
  8PyAirbyte supports structured JSON logging, which is disabled by default. To enable structured
  9logging in JSON, set `AIRBYTE_STRUCTURED_LOGGING` to `True`.
 10"""
 11
 12from __future__ import annotations
 13
 14import logging
 15import os
 16import platform
 17import sys
 18import tempfile
 19import warnings
 20from functools import lru_cache
 21from pathlib import Path
 22
 23import structlog
 24import ulid
 25
 26from airbyte_cdk.utils.datetime_helpers import ab_datetime_now
 27
 28from airbyte.constants import _str_to_bool
 29
 30
 31AIRBYTE_STRUCTURED_LOGGING: bool = _str_to_bool(
 32    os.getenv(key="AIRBYTE_STRUCTURED_LOGGING"),
 33    default=False,
 34)
 35"""Whether to enable structured logging.
 36
 37This value is read from the `AIRBYTE_STRUCTURED_LOGGING` environment variable. If the variable is
 38not set, the default value is `False`.
 39"""
 40
 41_warned_messages: set[str] = set()
 42
 43
 44def warn_once(
 45    message: str,
 46    logger: logging.Logger | None = None,
 47    *,
 48    with_stack: int | bool,
 49) -> None:
 50    """Emit a warning message only once.
 51
 52    This function is a wrapper around the `warnings.warn` function that logs the warning message
 53    to the global logger. The warning message is only emitted once per unique message.
 54    """
 55    if message in _warned_messages:
 56        return
 57
 58    if not with_stack:
 59        stacklevel = 0
 60    elif with_stack is True:
 61        stacklevel = 2
 62    elif isinstance(with_stack, int):
 63        stacklevel = with_stack
 64    else:
 65        stacklevel = 0
 66
 67    _warned_messages.add(message)
 68    warnings.warn(
 69        message,
 70        category=UserWarning,
 71        stacklevel=stacklevel,
 72    )
 73
 74    if logger:
 75        logger.warning(message)
 76
 77
 78def _get_logging_root() -> Path | None:
 79    """Return the root directory for logs.
 80
 81    Returns `None` if no valid path can be found.
 82
 83    This is the directory where logs are stored.
 84    """
 85    if "AIRBYTE_LOGGING_ROOT" in os.environ:
 86        log_root = Path(os.environ["AIRBYTE_LOGGING_ROOT"])
 87    elif platform.system() == "Darwin" or platform.system() == "Linux":
 88        # Use /tmp on macOS and Linux
 89        log_root = Path("/tmp") / "airbyte" / "logs"
 90    else:
 91        # Use the default temp directory on Windows or any other OS
 92        log_root = Path(tempfile.gettempdir()) / "airbyte" / "logs"
 93
 94    try:
 95        # Attempt to create the log root directory if it does not exist
 96        log_root.mkdir(parents=True, exist_ok=True)
 97    except OSError:
 98        # Handle the error by returning None
 99        warn_once(
100            (
101                f"Failed to create PyAirbyte logging directory at `{log_root}`. "
102                "You can override the default path by setting the `AIRBYTE_LOGGING_ROOT` "
103                "environment variable."
104            ),
105            with_stack=False,
106        )
107        return None
108    else:
109        return log_root
110
111
112AIRBYTE_LOGGING_ROOT: Path | None = _get_logging_root()
113"""The root directory for Airbyte logs.
114
115This value can be overridden by setting the `AIRBYTE_LOGGING_ROOT` environment variable.
116
117If not provided, PyAirbyte will use `/tmp/airbyte/logs/` where `/tmp/` is the OS's default
118temporary directory. If the directory cannot be created, PyAirbyte will log a warning and
119set this value to `None`.
120"""
121
122
123@lru_cache
124def get_global_file_logger() -> logging.Logger | None:
125    """Return the global logger for PyAirbyte.
126
127    This logger is configured to write logs to the console and to a file in the log directory.
128    """
129    logger = logging.getLogger("airbyte")
130    logger.setLevel(logging.INFO)
131    logger.propagate = False
132
133    if AIRBYTE_LOGGING_ROOT is None:
134        # No temp directory available, so return None
135        return None
136
137    # Else, configure the logger to write to a file
138
139    # Remove any existing handlers
140    for handler in logger.handlers:
141        logger.removeHandler(handler)
142
143    yyyy_mm_dd: str = ab_datetime_now().strftime("%Y-%m-%d")
144    folder = AIRBYTE_LOGGING_ROOT / yyyy_mm_dd
145    try:
146        folder.mkdir(parents=True, exist_ok=True)
147    except Exception:
148        warn_once(
149            f"Failed to create logging directory at '{folder!s}'.",
150            with_stack=False,
151        )
152        return None
153
154    logfile_path = folder / f"airbyte-log-{str(ulid.ULID())[2:11]}.log"
155    print(f"Writing PyAirbyte logs to file: {logfile_path!s}", file=sys.stderr)
156
157    file_handler = logging.FileHandler(
158        filename=logfile_path,
159        encoding="utf-8",
160    )
161
162    if AIRBYTE_STRUCTURED_LOGGING:
163        # Create a formatter and set it for the handler
164        formatter = logging.Formatter("%(message)s")
165        file_handler.setFormatter(formatter)
166
167        # Add the file handler to the logger
168        logger.addHandler(file_handler)
169
170        # Configure structlog
171        structlog.configure(
172            processors=[
173                structlog.processors.TimeStamper(fmt="%Y-%m-%d %H:%M:%S"),
174                structlog.stdlib.add_log_level,
175                structlog.stdlib.PositionalArgumentsFormatter(),
176                structlog.processors.StackInfoRenderer(),
177                structlog.processors.format_exc_info,
178                structlog.processors.JSONRenderer(),
179            ],
180            context_class=dict,
181            logger_factory=structlog.stdlib.LoggerFactory(),
182            wrapper_class=structlog.stdlib.BoundLogger,
183            cache_logger_on_first_use=True,
184        )
185
186        # Create a logger
187        return structlog.get_logger("airbyte")
188
189    # Create and configure file handler
190    file_handler.setFormatter(
191        logging.Formatter(
192            fmt="%(asctime)s - %(levelname)s - %(message)s",
193            datefmt="%Y-%m-%d %H:%M:%S",
194        )
195    )
196
197    logger.addHandler(file_handler)
198    return logger
199
200
201def get_global_stats_log_path() -> Path | None:
202    """Return the path to the performance log file."""
203    if AIRBYTE_LOGGING_ROOT is None:
204        return None
205
206    folder = AIRBYTE_LOGGING_ROOT
207    try:
208        folder.mkdir(parents=True, exist_ok=True)
209    except Exception:
210        warn_once(
211            f"Failed to create logging directory at '{folder!s}'.",
212            with_stack=False,
213        )
214        return None
215
216    return folder / "airbyte-stats.log"
217
218
219@lru_cache
220def get_global_stats_logger() -> structlog.BoundLogger:
221    """Create a stats logger for performance metrics."""
222    logger = logging.getLogger("airbyte.stats")
223    logger.setLevel(logging.INFO)
224    logger.propagate = False
225
226    # Configure structlog
227    structlog.configure(
228        processors=[
229            structlog.processors.TimeStamper(fmt="%Y-%m-%d %H:%M:%S"),
230            structlog.stdlib.PositionalArgumentsFormatter(),
231            structlog.processors.JSONRenderer(),
232        ],
233        context_class=dict,
234        logger_factory=structlog.stdlib.LoggerFactory(),
235        wrapper_class=structlog.stdlib.BoundLogger,
236        cache_logger_on_first_use=True,
237    )
238
239    logfile_path: Path | None = get_global_stats_log_path()
240    if AIRBYTE_LOGGING_ROOT is None or logfile_path is None:
241        # No temp directory available, so return no-op logger without handlers
242        return structlog.get_logger("airbyte.stats")
243
244    print(f"Writing PyAirbyte performance stats to file: {logfile_path!s}", file=sys.stderr)
245
246    # Remove any existing handlers
247    for handler in logger.handlers:
248        logger.removeHandler(handler)
249
250    folder = AIRBYTE_LOGGING_ROOT
251    try:
252        folder.mkdir(parents=True, exist_ok=True)
253    except Exception:
254        warn_once(
255            f"Failed to create logging directory at '{folder!s}'.",
256            with_stack=False,
257        )
258        return structlog.get_logger("airbyte.stats")
259
260    file_handler = logging.FileHandler(
261        filename=logfile_path,
262        encoding="utf-8",
263    )
264
265    # Create a formatter and set it for the handler
266    formatter = logging.Formatter("%(message)s")
267    file_handler.setFormatter(formatter)
268
269    # Add the file handler to the logger
270    logger.addHandler(file_handler)
271
272    # Create a logger
273    return structlog.get_logger("airbyte.stats")
274
275
276def new_passthrough_file_logger(connector_name: str) -> logging.Logger:
277    """Create a logger from logging module."""
278    logger = logging.getLogger(f"airbyte.{connector_name}")
279    logger.setLevel(logging.INFO)
280
281    # Prevent logging to stderr by stopping propagation to the root logger
282    logger.propagate = False
283
284    if AIRBYTE_LOGGING_ROOT is None:
285        # No temp directory available, so return a basic logger
286        return logger
287
288    # Else, configure the logger to write to a file
289
290    # Remove any existing handlers
291    for handler in logger.handlers:
292        logger.removeHandler(handler)
293
294    folder = AIRBYTE_LOGGING_ROOT / connector_name
295    folder.mkdir(parents=True, exist_ok=True)
296
297    # Create a file handler
298    global_logger = get_global_file_logger()
299    logfile_path = folder / f"{connector_name}-log-{str(ulid.ULID())[2:11]}.log"
300    logfile_msg = f"Writing `{connector_name}` logs to file: {logfile_path!s}"
301    print(logfile_msg, file=sys.stderr)
302    if global_logger:
303        global_logger.info(logfile_msg)
304
305    file_handler = logging.FileHandler(logfile_path)
306    file_handler.setLevel(logging.INFO)
307
308    if AIRBYTE_STRUCTURED_LOGGING:
309        # Create a formatter and set it for the handler
310        formatter = logging.Formatter("%(message)s")
311        file_handler.setFormatter(formatter)
312
313        # Add the file handler to the logger
314        logger.addHandler(file_handler)
315
316        # Configure structlog
317        structlog.configure(
318            processors=[
319                structlog.processors.TimeStamper(fmt="%Y-%m-%d %H:%M:%S"),
320                structlog.stdlib.add_log_level,
321                structlog.stdlib.PositionalArgumentsFormatter(),
322                structlog.processors.StackInfoRenderer(),
323                structlog.processors.format_exc_info,
324                structlog.processors.JSONRenderer(),
325            ],
326            context_class=dict,
327            logger_factory=structlog.stdlib.LoggerFactory(),
328            wrapper_class=structlog.stdlib.BoundLogger,
329            cache_logger_on_first_use=True,
330        )
331
332        # Create a logger
333        return structlog.get_logger(f"airbyte.{connector_name}")
334
335    # Else, write logs in plain text
336
337    file_handler.setFormatter(
338        logging.Formatter(
339            fmt="%(asctime)s - %(levelname)s - %(message)s",
340            datefmt="%Y-%m-%d %H:%M:%S",
341        )
342    )
343
344    logger.addHandler(file_handler)
345    return logger
AIRBYTE_STRUCTURED_LOGGING: bool = False

Whether to enable structured logging.

This value is read from the AIRBYTE_STRUCTURED_LOGGING environment variable. If the variable is not set, the default value is False.

def warn_once( message: str, logger: logging.Logger | None = None, *, with_stack: int | bool) -> None:
45def warn_once(
46    message: str,
47    logger: logging.Logger | None = None,
48    *,
49    with_stack: int | bool,
50) -> None:
51    """Emit a warning message only once.
52
53    This function is a wrapper around the `warnings.warn` function that logs the warning message
54    to the global logger. The warning message is only emitted once per unique message.
55    """
56    if message in _warned_messages:
57        return
58
59    if not with_stack:
60        stacklevel = 0
61    elif with_stack is True:
62        stacklevel = 2
63    elif isinstance(with_stack, int):
64        stacklevel = with_stack
65    else:
66        stacklevel = 0
67
68    _warned_messages.add(message)
69    warnings.warn(
70        message,
71        category=UserWarning,
72        stacklevel=stacklevel,
73    )
74
75    if logger:
76        logger.warning(message)

Emit a warning message only once.

This function is a wrapper around the warnings.warn function that logs the warning message to the global logger. The warning message is only emitted once per unique message.

AIRBYTE_LOGGING_ROOT: pathlib.Path | None = PosixPath('/tmp/airbyte/logs')

The root directory for Airbyte logs.

This value can be overridden by setting the AIRBYTE_LOGGING_ROOT environment variable.

If not provided, PyAirbyte will use /tmp/airbyte/logs/ where /tmp/ is the OS's default temporary directory. If the directory cannot be created, PyAirbyte will log a warning and set this value to None.

@lru_cache
def get_global_file_logger() -> logging.Logger | None:
124@lru_cache
125def get_global_file_logger() -> logging.Logger | None:
126    """Return the global logger for PyAirbyte.
127
128    This logger is configured to write logs to the console and to a file in the log directory.
129    """
130    logger = logging.getLogger("airbyte")
131    logger.setLevel(logging.INFO)
132    logger.propagate = False
133
134    if AIRBYTE_LOGGING_ROOT is None:
135        # No temp directory available, so return None
136        return None
137
138    # Else, configure the logger to write to a file
139
140    # Remove any existing handlers
141    for handler in logger.handlers:
142        logger.removeHandler(handler)
143
144    yyyy_mm_dd: str = ab_datetime_now().strftime("%Y-%m-%d")
145    folder = AIRBYTE_LOGGING_ROOT / yyyy_mm_dd
146    try:
147        folder.mkdir(parents=True, exist_ok=True)
148    except Exception:
149        warn_once(
150            f"Failed to create logging directory at '{folder!s}'.",
151            with_stack=False,
152        )
153        return None
154
155    logfile_path = folder / f"airbyte-log-{str(ulid.ULID())[2:11]}.log"
156    print(f"Writing PyAirbyte logs to file: {logfile_path!s}", file=sys.stderr)
157
158    file_handler = logging.FileHandler(
159        filename=logfile_path,
160        encoding="utf-8",
161    )
162
163    if AIRBYTE_STRUCTURED_LOGGING:
164        # Create a formatter and set it for the handler
165        formatter = logging.Formatter("%(message)s")
166        file_handler.setFormatter(formatter)
167
168        # Add the file handler to the logger
169        logger.addHandler(file_handler)
170
171        # Configure structlog
172        structlog.configure(
173            processors=[
174                structlog.processors.TimeStamper(fmt="%Y-%m-%d %H:%M:%S"),
175                structlog.stdlib.add_log_level,
176                structlog.stdlib.PositionalArgumentsFormatter(),
177                structlog.processors.StackInfoRenderer(),
178                structlog.processors.format_exc_info,
179                structlog.processors.JSONRenderer(),
180            ],
181            context_class=dict,
182            logger_factory=structlog.stdlib.LoggerFactory(),
183            wrapper_class=structlog.stdlib.BoundLogger,
184            cache_logger_on_first_use=True,
185        )
186
187        # Create a logger
188        return structlog.get_logger("airbyte")
189
190    # Create and configure file handler
191    file_handler.setFormatter(
192        logging.Formatter(
193            fmt="%(asctime)s - %(levelname)s - %(message)s",
194            datefmt="%Y-%m-%d %H:%M:%S",
195        )
196    )
197
198    logger.addHandler(file_handler)
199    return logger

Return the global logger for PyAirbyte.

This logger is configured to write logs to the console and to a file in the log directory.

def get_global_stats_log_path() -> pathlib.Path | None:
202def get_global_stats_log_path() -> Path | None:
203    """Return the path to the performance log file."""
204    if AIRBYTE_LOGGING_ROOT is None:
205        return None
206
207    folder = AIRBYTE_LOGGING_ROOT
208    try:
209        folder.mkdir(parents=True, exist_ok=True)
210    except Exception:
211        warn_once(
212            f"Failed to create logging directory at '{folder!s}'.",
213            with_stack=False,
214        )
215        return None
216
217    return folder / "airbyte-stats.log"

Return the path to the performance log file.

@lru_cache
def get_global_stats_logger() -> structlog._generic.BoundLogger:
220@lru_cache
221def get_global_stats_logger() -> structlog.BoundLogger:
222    """Create a stats logger for performance metrics."""
223    logger = logging.getLogger("airbyte.stats")
224    logger.setLevel(logging.INFO)
225    logger.propagate = False
226
227    # Configure structlog
228    structlog.configure(
229        processors=[
230            structlog.processors.TimeStamper(fmt="%Y-%m-%d %H:%M:%S"),
231            structlog.stdlib.PositionalArgumentsFormatter(),
232            structlog.processors.JSONRenderer(),
233        ],
234        context_class=dict,
235        logger_factory=structlog.stdlib.LoggerFactory(),
236        wrapper_class=structlog.stdlib.BoundLogger,
237        cache_logger_on_first_use=True,
238    )
239
240    logfile_path: Path | None = get_global_stats_log_path()
241    if AIRBYTE_LOGGING_ROOT is None or logfile_path is None:
242        # No temp directory available, so return no-op logger without handlers
243        return structlog.get_logger("airbyte.stats")
244
245    print(f"Writing PyAirbyte performance stats to file: {logfile_path!s}", file=sys.stderr)
246
247    # Remove any existing handlers
248    for handler in logger.handlers:
249        logger.removeHandler(handler)
250
251    folder = AIRBYTE_LOGGING_ROOT
252    try:
253        folder.mkdir(parents=True, exist_ok=True)
254    except Exception:
255        warn_once(
256            f"Failed to create logging directory at '{folder!s}'.",
257            with_stack=False,
258        )
259        return structlog.get_logger("airbyte.stats")
260
261    file_handler = logging.FileHandler(
262        filename=logfile_path,
263        encoding="utf-8",
264    )
265
266    # Create a formatter and set it for the handler
267    formatter = logging.Formatter("%(message)s")
268    file_handler.setFormatter(formatter)
269
270    # Add the file handler to the logger
271    logger.addHandler(file_handler)
272
273    # Create a logger
274    return structlog.get_logger("airbyte.stats")

Create a stats logger for performance metrics.

def new_passthrough_file_logger(connector_name: str) -> logging.Logger:
277def new_passthrough_file_logger(connector_name: str) -> logging.Logger:
278    """Create a logger from logging module."""
279    logger = logging.getLogger(f"airbyte.{connector_name}")
280    logger.setLevel(logging.INFO)
281
282    # Prevent logging to stderr by stopping propagation to the root logger
283    logger.propagate = False
284
285    if AIRBYTE_LOGGING_ROOT is None:
286        # No temp directory available, so return a basic logger
287        return logger
288
289    # Else, configure the logger to write to a file
290
291    # Remove any existing handlers
292    for handler in logger.handlers:
293        logger.removeHandler(handler)
294
295    folder = AIRBYTE_LOGGING_ROOT / connector_name
296    folder.mkdir(parents=True, exist_ok=True)
297
298    # Create a file handler
299    global_logger = get_global_file_logger()
300    logfile_path = folder / f"{connector_name}-log-{str(ulid.ULID())[2:11]}.log"
301    logfile_msg = f"Writing `{connector_name}` logs to file: {logfile_path!s}"
302    print(logfile_msg, file=sys.stderr)
303    if global_logger:
304        global_logger.info(logfile_msg)
305
306    file_handler = logging.FileHandler(logfile_path)
307    file_handler.setLevel(logging.INFO)
308
309    if AIRBYTE_STRUCTURED_LOGGING:
310        # Create a formatter and set it for the handler
311        formatter = logging.Formatter("%(message)s")
312        file_handler.setFormatter(formatter)
313
314        # Add the file handler to the logger
315        logger.addHandler(file_handler)
316
317        # Configure structlog
318        structlog.configure(
319            processors=[
320                structlog.processors.TimeStamper(fmt="%Y-%m-%d %H:%M:%S"),
321                structlog.stdlib.add_log_level,
322                structlog.stdlib.PositionalArgumentsFormatter(),
323                structlog.processors.StackInfoRenderer(),
324                structlog.processors.format_exc_info,
325                structlog.processors.JSONRenderer(),
326            ],
327            context_class=dict,
328            logger_factory=structlog.stdlib.LoggerFactory(),
329            wrapper_class=structlog.stdlib.BoundLogger,
330            cache_logger_on_first_use=True,
331        )
332
333        # Create a logger
334        return structlog.get_logger(f"airbyte.{connector_name}")
335
336    # Else, write logs in plain text
337
338    file_handler.setFormatter(
339        logging.Formatter(
340            fmt="%(asctime)s - %(levelname)s - %(message)s",
341            datefmt="%Y-%m-%d %H:%M:%S",
342        )
343    )
344
345    logger.addHandler(file_handler)
346    return logger

Create a logger from logging module.