airbyte.cloud.sync_results
Sync results for Airbyte Cloud workspaces.
Examples
Run a sync job and wait for completion
To get started, we'll need a .CloudConnection object. You can obtain this object by calling
.CloudWorkspace.get_connection().
from airbyte import cloud
# Initialize an Airbyte Cloud workspace object
workspace = cloud.CloudWorkspace(
workspace_id="123",
api_key=ab.get_secret("AIRBYTE_CLOUD_API_KEY"),
)
# Get a connection object
connection = workspace.get_connection(connection_id="456")
Once we have a .CloudConnection object, we can simply call run_sync()
to start a sync job and wait for it to complete.
# Run a sync job
sync_result: SyncResult = connection.run_sync()
Run a sync job and return immediately
By default, run_sync() will wait for the job to complete and raise an
exception if the job fails. You can instead return immediately by setting
wait=False.
# Start the sync job and return immediately
sync_result: SyncResult = connection.run_sync(wait=False)
while not sync_result.is_job_complete():
print("Job is still running...")
time.sleep(5)
print(f"Job is complete! Status: {sync_result.get_job_status()}")
Examining the sync result
You can examine the sync result to get more information about the job:
sync_result: SyncResult = connection.run_sync()
# Print the job details
print(
f'''
Job ID: {sync_result.job_id}
Job URL: {sync_result.job_url}
Start Time: {sync_result.start_time}
Records Synced: {sync_result.records_synced}
Bytes Synced: {sync_result.bytes_synced}
Job Status: {sync_result.get_job_status()}
List of Stream Names: {', '.join(sync_result.stream_names)}
'''
)
Reading data from Airbyte Cloud sync result
This feature is currently only available for specific SQL-based destinations. This includes
SQL-based destinations such as Snowflake and BigQuery. The list of supported destinations may be
determined by inspecting the constant airbyte.cloud.constants.READABLE_DESTINATION_TYPES.
If your destination is supported, you can read records directly from the SyncResult object.
# Assuming we've already created a `connection` object...
sync_result = connection.get_sync_result()
# Print a list of available stream names
print(sync_result.stream_names)
# Get a dataset from the sync result
dataset: CachedDataset = sync_result.get_dataset("users")
# Get the SQLAlchemy table to use in SQL queries...
users_table = dataset.to_sql_table()
print(f"Table name: {users_table.name}")
# Or iterate over the dataset directly
for record in dataset:
print(record)
1# Copyright (c) 2024 Airbyte, Inc., all rights reserved. 2"""Sync results for Airbyte Cloud workspaces. 3 4## Examples 5 6### Run a sync job and wait for completion 7 8To get started, we'll need a `.CloudConnection` object. You can obtain this object by calling 9`.CloudWorkspace.get_connection()`. 10 11```python 12from airbyte import cloud 13 14# Initialize an Airbyte Cloud workspace object 15workspace = cloud.CloudWorkspace( 16 workspace_id="123", 17 api_key=ab.get_secret("AIRBYTE_CLOUD_API_KEY"), 18) 19 20# Get a connection object 21connection = workspace.get_connection(connection_id="456") 22``` 23 24Once we have a `.CloudConnection` object, we can simply call `run_sync()` 25to start a sync job and wait for it to complete. 26 27```python 28# Run a sync job 29sync_result: SyncResult = connection.run_sync() 30``` 31 32### Run a sync job and return immediately 33 34By default, `run_sync()` will wait for the job to complete and raise an 35exception if the job fails. You can instead return immediately by setting 36`wait=False`. 37 38```python 39# Start the sync job and return immediately 40sync_result: SyncResult = connection.run_sync(wait=False) 41 42while not sync_result.is_job_complete(): 43 print("Job is still running...") 44 time.sleep(5) 45 46print(f"Job is complete! Status: {sync_result.get_job_status()}") 47``` 48 49### Examining the sync result 50 51You can examine the sync result to get more information about the job: 52 53```python 54sync_result: SyncResult = connection.run_sync() 55 56# Print the job details 57print( 58 f''' 59 Job ID: {sync_result.job_id} 60 Job URL: {sync_result.job_url} 61 Start Time: {sync_result.start_time} 62 Records Synced: {sync_result.records_synced} 63 Bytes Synced: {sync_result.bytes_synced} 64 Job Status: {sync_result.get_job_status()} 65 List of Stream Names: {', '.join(sync_result.stream_names)} 66 ''' 67) 68``` 69 70### Reading data from Airbyte Cloud sync result 71 72**This feature is currently only available for specific SQL-based destinations.** This includes 73SQL-based destinations such as Snowflake and BigQuery. The list of supported destinations may be 74determined by inspecting the constant `airbyte.cloud.constants.READABLE_DESTINATION_TYPES`. 75 76If your destination is supported, you can read records directly from the SyncResult object. 77 78```python 79# Assuming we've already created a `connection` object... 80sync_result = connection.get_sync_result() 81 82# Print a list of available stream names 83print(sync_result.stream_names) 84 85# Get a dataset from the sync result 86dataset: CachedDataset = sync_result.get_dataset("users") 87 88# Get the SQLAlchemy table to use in SQL queries... 89users_table = dataset.to_sql_table() 90print(f"Table name: {users_table.name}") 91 92# Or iterate over the dataset directly 93for record in dataset: 94 print(record) 95``` 96 97------ 98 99""" 100 101from __future__ import annotations 102 103import time 104from collections.abc import Iterator, Mapping 105from dataclasses import asdict, dataclass 106from typing import TYPE_CHECKING, Any 107 108from typing_extensions import final 109 110from airbyte_cdk.utils.datetime_helpers import ab_datetime_parse 111 112from airbyte._util import api_util 113from airbyte.caches._utils._dest_to_cache import destination_to_cache 114from airbyte.cloud.constants import FAILED_STATUSES, FINAL_STATUSES 115from airbyte.cloud.models import CloudConnectionInfo, CloudJobInfo, JobStatusEnum 116from airbyte.datasets import CachedDataset 117from airbyte.exceptions import AirbyteConnectionSyncError, AirbyteConnectionSyncTimeoutError 118 119 120DEFAULT_SYNC_TIMEOUT_SECONDS = 30 * 60 # 30 minutes 121"""The default timeout for waiting for a sync job to complete, in seconds.""" 122 123if TYPE_CHECKING: 124 from datetime import datetime 125 126 import sqlalchemy 127 128 from airbyte.caches.base import CacheBase 129 from airbyte.cloud.connections import CloudConnection 130 from airbyte.cloud.workspaces import CloudWorkspace 131 132 133@dataclass 134class SyncAttempt: 135 """Represents a single attempt of a sync job. 136 137 **This class is not meant to be instantiated directly.** Instead, obtain a `SyncAttempt` by 138 calling `.SyncResult.get_attempts()`. 139 """ 140 141 workspace: CloudWorkspace 142 connection: CloudConnection 143 job_id: int 144 attempt_number: int 145 _attempt_data: dict[str, Any] | None = None 146 147 @property 148 def attempt_id(self) -> int: 149 """Return the attempt ID.""" 150 return self._get_attempt_data()["id"] 151 152 @property 153 def status(self) -> str: 154 """Return the attempt status.""" 155 return self._get_attempt_data()["status"] 156 157 @property 158 def bytes_synced(self) -> int: 159 """Return the number of bytes synced in this attempt.""" 160 return self._get_attempt_data().get("bytesSynced", 0) 161 162 @property 163 def records_synced(self) -> int: 164 """Return the number of records synced in this attempt.""" 165 return self._get_attempt_data().get("recordsSynced", 0) 166 167 @property 168 def created_at(self) -> datetime: 169 """Return the creation time of the attempt.""" 170 timestamp = self._get_attempt_data()["createdAt"] 171 return ab_datetime_parse(timestamp) 172 173 def _get_attempt_data(self) -> dict[str, Any]: 174 """Get attempt data from the provided attempt data.""" 175 if self._attempt_data is None: 176 raise ValueError( 177 "Attempt data not provided. SyncAttempt should be created via " 178 "SyncResult.get_attempts()." 179 ) 180 return self._attempt_data["attempt"] 181 182 def get_full_log_text(self) -> str: 183 """Return the complete log text for this attempt. 184 185 Returns: 186 String containing all log text for this attempt, with lines separated by newlines. 187 """ 188 if self._attempt_data is None: 189 return "" 190 191 logs_data = self._attempt_data.get("logs") 192 if not logs_data: 193 return "" 194 195 result = "" 196 197 if "events" in logs_data: 198 log_events = logs_data["events"] 199 if log_events: 200 log_lines = [] 201 for event in log_events: 202 timestamp = event.get("timestamp", "") 203 level = event.get("level", "INFO") 204 message = event.get("message", "") 205 log_lines.append( 206 f"[{timestamp}] {level}: {message}" # pyrefly: ignore[bad-argument-type] 207 ) 208 result = "\n".join(log_lines) 209 elif "logLines" in logs_data: 210 log_lines = logs_data["logLines"] 211 if log_lines: 212 result = "\n".join(log_lines) 213 214 return result 215 216 217@dataclass 218class SyncResult: 219 """The result of a sync operation. 220 221 **This class is not meant to be instantiated directly.** Instead, obtain a `SyncResult` by 222 interacting with the `.CloudWorkspace` and `.CloudConnection` objects. 223 """ 224 225 workspace: CloudWorkspace 226 connection: CloudConnection 227 job_id: int 228 table_name_prefix: str = "" 229 table_name_suffix: str = "" 230 _latest_job_info: CloudJobInfo | None = None 231 _connection_response: CloudConnectionInfo | None = None 232 _cache: CacheBase | None = None 233 _job_with_attempts_info: dict[str, Any] | None = None 234 235 @property 236 def job_url(self) -> str: 237 """Return the URL of the sync job. 238 239 Note: This currently returns the connection's job history URL, as there is no direct URL 240 to a specific job in the Airbyte Cloud web app. 241 242 TODO: Implement a direct job logs URL on top of the event-id of the specific attempt number. 243 E.g. {self.connection.job_history_url}?eventId={event-guid}&openLogs=true 244 """ 245 return f"{self.connection.job_history_url}" 246 247 def _get_connection_info(self, *, force_refresh: bool = False) -> CloudConnectionInfo: 248 """Return connection info for the sync job.""" 249 if self._connection_response and not force_refresh: 250 return self._connection_response 251 252 self._connection_response = CloudConnectionInfo.from_api_response( 253 api_util.get_connection( 254 workspace_id=self.workspace.workspace_id, 255 api_root=self.workspace.api_root, 256 connection_id=self.connection.connection_id, 257 client_id=self.workspace.client_id, 258 client_secret=self.workspace.client_secret, 259 bearer_token=self.workspace.bearer_token, 260 ) 261 ) 262 return self._connection_response 263 264 def _get_destination_configuration(self, *, force_refresh: bool = False) -> dict[str, Any]: 265 """Return the destination configuration for the sync job.""" 266 connection_info = self._get_connection_info(force_refresh=force_refresh) 267 destination_response = api_util.get_destination( 268 destination_id=connection_info.destination_id, 269 api_root=self.workspace.api_root, 270 client_id=self.workspace.client_id, 271 client_secret=self.workspace.client_secret, 272 bearer_token=self.workspace.bearer_token, 273 ) 274 configuration = destination_response.configuration 275 if isinstance(configuration, Mapping): 276 configuration_dict = configuration 277 else: 278 configuration_dict = asdict(configuration) 279 return {**configuration_dict, "destinationType": destination_response.destination_type} 280 281 def is_job_complete(self) -> bool: 282 """Check if the sync job is complete.""" 283 return self.get_job_status() in FINAL_STATUSES 284 285 def get_job_status(self) -> JobStatusEnum: 286 """Check if the sync job is still running.""" 287 return self._fetch_latest_job_info().status 288 289 def _fetch_latest_job_info(self) -> CloudJobInfo: 290 """Return the job info for the sync job.""" 291 if self._latest_job_info and self._latest_job_info.status in FINAL_STATUSES: 292 return self._latest_job_info 293 294 self._latest_job_info = CloudJobInfo.from_api_response( 295 api_util.get_job_info( 296 job_id=self.job_id, 297 api_root=self.workspace.api_root, 298 client_id=self.workspace.client_id, 299 client_secret=self.workspace.client_secret, 300 bearer_token=self.workspace.bearer_token, 301 ) 302 ) 303 return self._latest_job_info 304 305 @property 306 def bytes_synced(self) -> int: 307 """Return the number of records processed.""" 308 return self._fetch_latest_job_info().bytes_synced or 0 309 310 @property 311 def records_synced(self) -> int: 312 """Return the number of records processed.""" 313 return self._fetch_latest_job_info().rows_synced or 0 314 315 @property 316 def start_time(self) -> datetime: 317 """Return the start time of the sync job in UTC.""" 318 try: 319 return ab_datetime_parse(self._fetch_latest_job_info().start_time) 320 except (ValueError, TypeError) as e: 321 if "Invalid isoformat string" in str(e): 322 job_info_raw = api_util._make_config_api_request( # noqa: SLF001 323 api_root=self.workspace.api_root, 324 config_api_root=self.workspace.config_api_root, 325 path="/jobs/get", 326 json={"id": self.job_id}, 327 client_id=self.workspace.client_id, 328 client_secret=self.workspace.client_secret, 329 bearer_token=self.workspace.bearer_token, 330 ) 331 raw_start_time = job_info_raw.get("startTime") 332 if raw_start_time: 333 return ab_datetime_parse(raw_start_time) 334 raise 335 336 def _fetch_job_with_attempts(self) -> dict[str, Any]: 337 """Fetch job info with attempts from Config API using lazy loading pattern.""" 338 if self._job_with_attempts_info is not None: 339 return self._job_with_attempts_info 340 341 self._job_with_attempts_info = api_util._make_config_api_request( # noqa: SLF001 # Config API helper 342 api_root=self.workspace.api_root, 343 config_api_root=self.workspace.config_api_root, 344 path="/jobs/get", 345 json={ 346 "id": self.job_id, 347 }, 348 client_id=self.workspace.client_id, 349 client_secret=self.workspace.client_secret, 350 bearer_token=self.workspace.bearer_token, 351 ) 352 return self._job_with_attempts_info 353 354 def get_attempts(self) -> list[SyncAttempt]: 355 """Return a list of attempts for this sync job.""" 356 job_with_attempts = self._fetch_job_with_attempts() 357 attempts_data = job_with_attempts.get("attempts", []) 358 359 return [ 360 SyncAttempt( 361 workspace=self.workspace, 362 connection=self.connection, 363 job_id=self.job_id, 364 attempt_number=i, 365 _attempt_data=attempt_data, 366 ) 367 for i, attempt_data in enumerate(attempts_data, start=0) 368 ] 369 370 def raise_failure_status( 371 self, 372 *, 373 refresh_status: bool = False, 374 ) -> None: 375 """Raise an exception if the sync job failed. 376 377 By default, this method will use the latest status available. If you want to refresh the 378 status before checking for failure, set `refresh_status=True`. If the job has failed, this 379 method will raise a `AirbyteConnectionSyncError`. 380 381 Otherwise, do nothing. 382 """ 383 if not refresh_status and self._latest_job_info: 384 latest_status = self._latest_job_info.status 385 else: 386 latest_status = self.get_job_status() 387 388 if latest_status in FAILED_STATUSES: 389 raise AirbyteConnectionSyncError( 390 workspace=self.workspace, 391 connection_id=self.connection.connection_id, 392 job_id=self.job_id, 393 job_status=self.get_job_status(), 394 ) 395 396 def wait_for_completion( 397 self, 398 *, 399 wait_timeout: int = DEFAULT_SYNC_TIMEOUT_SECONDS, 400 raise_timeout: bool = True, 401 raise_failure: bool = False, 402 ) -> JobStatusEnum: 403 """Wait for a job to finish running.""" 404 start_time = time.time() 405 while True: 406 latest_status = self.get_job_status() 407 if latest_status in FINAL_STATUSES: 408 if raise_failure: 409 # No-op if the job succeeded or is still running: 410 self.raise_failure_status() 411 412 return latest_status 413 414 if time.time() - start_time > wait_timeout: 415 if raise_timeout: 416 raise AirbyteConnectionSyncTimeoutError( 417 workspace=self.workspace, 418 connection_id=self.connection.connection_id, 419 job_id=self.job_id, 420 job_status=latest_status, 421 timeout=wait_timeout, 422 ) 423 424 return latest_status # This will be a non-final status 425 426 time.sleep(api_util.JOB_WAIT_INTERVAL_SECS) 427 428 def get_sql_cache(self) -> CacheBase: 429 """Return a SQL Cache object for working with the data in a SQL-based destination's.""" 430 if self._cache: 431 return self._cache 432 433 destination_configuration = self._get_destination_configuration() 434 self._cache = destination_to_cache(destination_configuration=destination_configuration) 435 return self._cache 436 437 def get_sql_engine(self) -> sqlalchemy.engine.Engine: 438 """Return a SQL Engine for querying a SQL-based destination.""" 439 return self.get_sql_cache().get_sql_engine() 440 441 def get_sql_table_name(self, stream_name: str) -> str: 442 """Return the SQL table name of the named stream.""" 443 return self.get_sql_cache().processor.get_sql_table_name(stream_name=stream_name) 444 445 def get_sql_table( 446 self, 447 stream_name: str, 448 ) -> sqlalchemy.Table: 449 """Return a SQLAlchemy table object for the named stream.""" 450 return self.get_sql_cache().processor.get_sql_table(stream_name) 451 452 def get_dataset(self, stream_name: str) -> CachedDataset: 453 """Retrieve an `airbyte.datasets.CachedDataset` object for a given stream name. 454 455 This can be used to read and analyze the data in a SQL-based destination. 456 457 TODO: In a future iteration, we can consider providing stream configuration information 458 (catalog information) to the `CachedDataset` object via the "Get stream properties" 459 API: https://reference.airbyte.com/reference/getstreamproperties 460 """ 461 return CachedDataset( 462 self.get_sql_cache(), 463 stream_name=stream_name, 464 stream_configuration=False, # Don't look for stream configuration in cache. 465 ) 466 467 def get_sql_database_name(self) -> str: 468 """Return the SQL database name.""" 469 cache = self.get_sql_cache() 470 return cache.get_database_name() 471 472 def get_sql_schema_name(self) -> str: 473 """Return the SQL schema name.""" 474 cache = self.get_sql_cache() 475 return cache.schema_name 476 477 @property 478 def stream_names(self) -> list[str]: 479 """Return the set of stream names.""" 480 return self.connection.stream_names 481 482 @final 483 @property 484 def streams( 485 self, 486 ) -> _SyncResultStreams: # pyrefly: ignore[unknown-name] 487 """Return a mapping of stream names to `airbyte.CachedDataset` objects. 488 489 This is a convenience wrapper around the `stream_names` 490 property and `get_dataset()` method. 491 """ 492 return self._SyncResultStreams(self) 493 494 class _SyncResultStreams(Mapping[str, CachedDataset]): 495 """A mapping of stream names to cached datasets.""" 496 497 def __init__( 498 self, 499 parent: SyncResult, 500 /, 501 ) -> None: 502 self.parent: SyncResult = parent 503 504 def __getitem__(self, key: str) -> CachedDataset: 505 return self.parent.get_dataset(stream_name=key) 506 507 def __iter__(self) -> Iterator[str]: 508 return iter(self.parent.stream_names) 509 510 def __len__(self) -> int: 511 return len(self.parent.stream_names) 512 513 514__all__ = [ 515 "SyncResult", 516 "SyncAttempt", 517]
218@dataclass 219class SyncResult: 220 """The result of a sync operation. 221 222 **This class is not meant to be instantiated directly.** Instead, obtain a `SyncResult` by 223 interacting with the `.CloudWorkspace` and `.CloudConnection` objects. 224 """ 225 226 workspace: CloudWorkspace 227 connection: CloudConnection 228 job_id: int 229 table_name_prefix: str = "" 230 table_name_suffix: str = "" 231 _latest_job_info: CloudJobInfo | None = None 232 _connection_response: CloudConnectionInfo | None = None 233 _cache: CacheBase | None = None 234 _job_with_attempts_info: dict[str, Any] | None = None 235 236 @property 237 def job_url(self) -> str: 238 """Return the URL of the sync job. 239 240 Note: This currently returns the connection's job history URL, as there is no direct URL 241 to a specific job in the Airbyte Cloud web app. 242 243 TODO: Implement a direct job logs URL on top of the event-id of the specific attempt number. 244 E.g. {self.connection.job_history_url}?eventId={event-guid}&openLogs=true 245 """ 246 return f"{self.connection.job_history_url}" 247 248 def _get_connection_info(self, *, force_refresh: bool = False) -> CloudConnectionInfo: 249 """Return connection info for the sync job.""" 250 if self._connection_response and not force_refresh: 251 return self._connection_response 252 253 self._connection_response = CloudConnectionInfo.from_api_response( 254 api_util.get_connection( 255 workspace_id=self.workspace.workspace_id, 256 api_root=self.workspace.api_root, 257 connection_id=self.connection.connection_id, 258 client_id=self.workspace.client_id, 259 client_secret=self.workspace.client_secret, 260 bearer_token=self.workspace.bearer_token, 261 ) 262 ) 263 return self._connection_response 264 265 def _get_destination_configuration(self, *, force_refresh: bool = False) -> dict[str, Any]: 266 """Return the destination configuration for the sync job.""" 267 connection_info = self._get_connection_info(force_refresh=force_refresh) 268 destination_response = api_util.get_destination( 269 destination_id=connection_info.destination_id, 270 api_root=self.workspace.api_root, 271 client_id=self.workspace.client_id, 272 client_secret=self.workspace.client_secret, 273 bearer_token=self.workspace.bearer_token, 274 ) 275 configuration = destination_response.configuration 276 if isinstance(configuration, Mapping): 277 configuration_dict = configuration 278 else: 279 configuration_dict = asdict(configuration) 280 return {**configuration_dict, "destinationType": destination_response.destination_type} 281 282 def is_job_complete(self) -> bool: 283 """Check if the sync job is complete.""" 284 return self.get_job_status() in FINAL_STATUSES 285 286 def get_job_status(self) -> JobStatusEnum: 287 """Check if the sync job is still running.""" 288 return self._fetch_latest_job_info().status 289 290 def _fetch_latest_job_info(self) -> CloudJobInfo: 291 """Return the job info for the sync job.""" 292 if self._latest_job_info and self._latest_job_info.status in FINAL_STATUSES: 293 return self._latest_job_info 294 295 self._latest_job_info = CloudJobInfo.from_api_response( 296 api_util.get_job_info( 297 job_id=self.job_id, 298 api_root=self.workspace.api_root, 299 client_id=self.workspace.client_id, 300 client_secret=self.workspace.client_secret, 301 bearer_token=self.workspace.bearer_token, 302 ) 303 ) 304 return self._latest_job_info 305 306 @property 307 def bytes_synced(self) -> int: 308 """Return the number of records processed.""" 309 return self._fetch_latest_job_info().bytes_synced or 0 310 311 @property 312 def records_synced(self) -> int: 313 """Return the number of records processed.""" 314 return self._fetch_latest_job_info().rows_synced or 0 315 316 @property 317 def start_time(self) -> datetime: 318 """Return the start time of the sync job in UTC.""" 319 try: 320 return ab_datetime_parse(self._fetch_latest_job_info().start_time) 321 except (ValueError, TypeError) as e: 322 if "Invalid isoformat string" in str(e): 323 job_info_raw = api_util._make_config_api_request( # noqa: SLF001 324 api_root=self.workspace.api_root, 325 config_api_root=self.workspace.config_api_root, 326 path="/jobs/get", 327 json={"id": self.job_id}, 328 client_id=self.workspace.client_id, 329 client_secret=self.workspace.client_secret, 330 bearer_token=self.workspace.bearer_token, 331 ) 332 raw_start_time = job_info_raw.get("startTime") 333 if raw_start_time: 334 return ab_datetime_parse(raw_start_time) 335 raise 336 337 def _fetch_job_with_attempts(self) -> dict[str, Any]: 338 """Fetch job info with attempts from Config API using lazy loading pattern.""" 339 if self._job_with_attempts_info is not None: 340 return self._job_with_attempts_info 341 342 self._job_with_attempts_info = api_util._make_config_api_request( # noqa: SLF001 # Config API helper 343 api_root=self.workspace.api_root, 344 config_api_root=self.workspace.config_api_root, 345 path="/jobs/get", 346 json={ 347 "id": self.job_id, 348 }, 349 client_id=self.workspace.client_id, 350 client_secret=self.workspace.client_secret, 351 bearer_token=self.workspace.bearer_token, 352 ) 353 return self._job_with_attempts_info 354 355 def get_attempts(self) -> list[SyncAttempt]: 356 """Return a list of attempts for this sync job.""" 357 job_with_attempts = self._fetch_job_with_attempts() 358 attempts_data = job_with_attempts.get("attempts", []) 359 360 return [ 361 SyncAttempt( 362 workspace=self.workspace, 363 connection=self.connection, 364 job_id=self.job_id, 365 attempt_number=i, 366 _attempt_data=attempt_data, 367 ) 368 for i, attempt_data in enumerate(attempts_data, start=0) 369 ] 370 371 def raise_failure_status( 372 self, 373 *, 374 refresh_status: bool = False, 375 ) -> None: 376 """Raise an exception if the sync job failed. 377 378 By default, this method will use the latest status available. If you want to refresh the 379 status before checking for failure, set `refresh_status=True`. If the job has failed, this 380 method will raise a `AirbyteConnectionSyncError`. 381 382 Otherwise, do nothing. 383 """ 384 if not refresh_status and self._latest_job_info: 385 latest_status = self._latest_job_info.status 386 else: 387 latest_status = self.get_job_status() 388 389 if latest_status in FAILED_STATUSES: 390 raise AirbyteConnectionSyncError( 391 workspace=self.workspace, 392 connection_id=self.connection.connection_id, 393 job_id=self.job_id, 394 job_status=self.get_job_status(), 395 ) 396 397 def wait_for_completion( 398 self, 399 *, 400 wait_timeout: int = DEFAULT_SYNC_TIMEOUT_SECONDS, 401 raise_timeout: bool = True, 402 raise_failure: bool = False, 403 ) -> JobStatusEnum: 404 """Wait for a job to finish running.""" 405 start_time = time.time() 406 while True: 407 latest_status = self.get_job_status() 408 if latest_status in FINAL_STATUSES: 409 if raise_failure: 410 # No-op if the job succeeded or is still running: 411 self.raise_failure_status() 412 413 return latest_status 414 415 if time.time() - start_time > wait_timeout: 416 if raise_timeout: 417 raise AirbyteConnectionSyncTimeoutError( 418 workspace=self.workspace, 419 connection_id=self.connection.connection_id, 420 job_id=self.job_id, 421 job_status=latest_status, 422 timeout=wait_timeout, 423 ) 424 425 return latest_status # This will be a non-final status 426 427 time.sleep(api_util.JOB_WAIT_INTERVAL_SECS) 428 429 def get_sql_cache(self) -> CacheBase: 430 """Return a SQL Cache object for working with the data in a SQL-based destination's.""" 431 if self._cache: 432 return self._cache 433 434 destination_configuration = self._get_destination_configuration() 435 self._cache = destination_to_cache(destination_configuration=destination_configuration) 436 return self._cache 437 438 def get_sql_engine(self) -> sqlalchemy.engine.Engine: 439 """Return a SQL Engine for querying a SQL-based destination.""" 440 return self.get_sql_cache().get_sql_engine() 441 442 def get_sql_table_name(self, stream_name: str) -> str: 443 """Return the SQL table name of the named stream.""" 444 return self.get_sql_cache().processor.get_sql_table_name(stream_name=stream_name) 445 446 def get_sql_table( 447 self, 448 stream_name: str, 449 ) -> sqlalchemy.Table: 450 """Return a SQLAlchemy table object for the named stream.""" 451 return self.get_sql_cache().processor.get_sql_table(stream_name) 452 453 def get_dataset(self, stream_name: str) -> CachedDataset: 454 """Retrieve an `airbyte.datasets.CachedDataset` object for a given stream name. 455 456 This can be used to read and analyze the data in a SQL-based destination. 457 458 TODO: In a future iteration, we can consider providing stream configuration information 459 (catalog information) to the `CachedDataset` object via the "Get stream properties" 460 API: https://reference.airbyte.com/reference/getstreamproperties 461 """ 462 return CachedDataset( 463 self.get_sql_cache(), 464 stream_name=stream_name, 465 stream_configuration=False, # Don't look for stream configuration in cache. 466 ) 467 468 def get_sql_database_name(self) -> str: 469 """Return the SQL database name.""" 470 cache = self.get_sql_cache() 471 return cache.get_database_name() 472 473 def get_sql_schema_name(self) -> str: 474 """Return the SQL schema name.""" 475 cache = self.get_sql_cache() 476 return cache.schema_name 477 478 @property 479 def stream_names(self) -> list[str]: 480 """Return the set of stream names.""" 481 return self.connection.stream_names 482 483 @final 484 @property 485 def streams( 486 self, 487 ) -> _SyncResultStreams: # pyrefly: ignore[unknown-name] 488 """Return a mapping of stream names to `airbyte.CachedDataset` objects. 489 490 This is a convenience wrapper around the `stream_names` 491 property and `get_dataset()` method. 492 """ 493 return self._SyncResultStreams(self) 494 495 class _SyncResultStreams(Mapping[str, CachedDataset]): 496 """A mapping of stream names to cached datasets.""" 497 498 def __init__( 499 self, 500 parent: SyncResult, 501 /, 502 ) -> None: 503 self.parent: SyncResult = parent 504 505 def __getitem__(self, key: str) -> CachedDataset: 506 return self.parent.get_dataset(stream_name=key) 507 508 def __iter__(self) -> Iterator[str]: 509 return iter(self.parent.stream_names) 510 511 def __len__(self) -> int: 512 return len(self.parent.stream_names)
The result of a sync operation.
This class is not meant to be instantiated directly. Instead, obtain a SyncResult by
interacting with the .CloudWorkspace and .CloudConnection objects.
236 @property 237 def job_url(self) -> str: 238 """Return the URL of the sync job. 239 240 Note: This currently returns the connection's job history URL, as there is no direct URL 241 to a specific job in the Airbyte Cloud web app. 242 243 TODO: Implement a direct job logs URL on top of the event-id of the specific attempt number. 244 E.g. {self.connection.job_history_url}?eventId={event-guid}&openLogs=true 245 """ 246 return f"{self.connection.job_history_url}"
Return the URL of the sync job.
Note: This currently returns the connection's job history URL, as there is no direct URL to a specific job in the Airbyte Cloud web app.
TODO: Implement a direct job logs URL on top of the event-id of the specific attempt number. E.g. {self.connection.job_history_url}?eventId={event-guid}&openLogs=true
282 def is_job_complete(self) -> bool: 283 """Check if the sync job is complete.""" 284 return self.get_job_status() in FINAL_STATUSES
Check if the sync job is complete.
286 def get_job_status(self) -> JobStatusEnum: 287 """Check if the sync job is still running.""" 288 return self._fetch_latest_job_info().status
Check if the sync job is still running.
306 @property 307 def bytes_synced(self) -> int: 308 """Return the number of records processed.""" 309 return self._fetch_latest_job_info().bytes_synced or 0
Return the number of records processed.
311 @property 312 def records_synced(self) -> int: 313 """Return the number of records processed.""" 314 return self._fetch_latest_job_info().rows_synced or 0
Return the number of records processed.
316 @property 317 def start_time(self) -> datetime: 318 """Return the start time of the sync job in UTC.""" 319 try: 320 return ab_datetime_parse(self._fetch_latest_job_info().start_time) 321 except (ValueError, TypeError) as e: 322 if "Invalid isoformat string" in str(e): 323 job_info_raw = api_util._make_config_api_request( # noqa: SLF001 324 api_root=self.workspace.api_root, 325 config_api_root=self.workspace.config_api_root, 326 path="/jobs/get", 327 json={"id": self.job_id}, 328 client_id=self.workspace.client_id, 329 client_secret=self.workspace.client_secret, 330 bearer_token=self.workspace.bearer_token, 331 ) 332 raw_start_time = job_info_raw.get("startTime") 333 if raw_start_time: 334 return ab_datetime_parse(raw_start_time) 335 raise
Return the start time of the sync job in UTC.
355 def get_attempts(self) -> list[SyncAttempt]: 356 """Return a list of attempts for this sync job.""" 357 job_with_attempts = self._fetch_job_with_attempts() 358 attempts_data = job_with_attempts.get("attempts", []) 359 360 return [ 361 SyncAttempt( 362 workspace=self.workspace, 363 connection=self.connection, 364 job_id=self.job_id, 365 attempt_number=i, 366 _attempt_data=attempt_data, 367 ) 368 for i, attempt_data in enumerate(attempts_data, start=0) 369 ]
Return a list of attempts for this sync job.
371 def raise_failure_status( 372 self, 373 *, 374 refresh_status: bool = False, 375 ) -> None: 376 """Raise an exception if the sync job failed. 377 378 By default, this method will use the latest status available. If you want to refresh the 379 status before checking for failure, set `refresh_status=True`. If the job has failed, this 380 method will raise a `AirbyteConnectionSyncError`. 381 382 Otherwise, do nothing. 383 """ 384 if not refresh_status and self._latest_job_info: 385 latest_status = self._latest_job_info.status 386 else: 387 latest_status = self.get_job_status() 388 389 if latest_status in FAILED_STATUSES: 390 raise AirbyteConnectionSyncError( 391 workspace=self.workspace, 392 connection_id=self.connection.connection_id, 393 job_id=self.job_id, 394 job_status=self.get_job_status(), 395 )
Raise an exception if the sync job failed.
By default, this method will use the latest status available. If you want to refresh the
status before checking for failure, set refresh_status=True. If the job has failed, this
method will raise a AirbyteConnectionSyncError.
Otherwise, do nothing.
397 def wait_for_completion( 398 self, 399 *, 400 wait_timeout: int = DEFAULT_SYNC_TIMEOUT_SECONDS, 401 raise_timeout: bool = True, 402 raise_failure: bool = False, 403 ) -> JobStatusEnum: 404 """Wait for a job to finish running.""" 405 start_time = time.time() 406 while True: 407 latest_status = self.get_job_status() 408 if latest_status in FINAL_STATUSES: 409 if raise_failure: 410 # No-op if the job succeeded or is still running: 411 self.raise_failure_status() 412 413 return latest_status 414 415 if time.time() - start_time > wait_timeout: 416 if raise_timeout: 417 raise AirbyteConnectionSyncTimeoutError( 418 workspace=self.workspace, 419 connection_id=self.connection.connection_id, 420 job_id=self.job_id, 421 job_status=latest_status, 422 timeout=wait_timeout, 423 ) 424 425 return latest_status # This will be a non-final status 426 427 time.sleep(api_util.JOB_WAIT_INTERVAL_SECS)
Wait for a job to finish running.
429 def get_sql_cache(self) -> CacheBase: 430 """Return a SQL Cache object for working with the data in a SQL-based destination's.""" 431 if self._cache: 432 return self._cache 433 434 destination_configuration = self._get_destination_configuration() 435 self._cache = destination_to_cache(destination_configuration=destination_configuration) 436 return self._cache
Return a SQL Cache object for working with the data in a SQL-based destination's.
438 def get_sql_engine(self) -> sqlalchemy.engine.Engine: 439 """Return a SQL Engine for querying a SQL-based destination.""" 440 return self.get_sql_cache().get_sql_engine()
Return a SQL Engine for querying a SQL-based destination.
442 def get_sql_table_name(self, stream_name: str) -> str: 443 """Return the SQL table name of the named stream.""" 444 return self.get_sql_cache().processor.get_sql_table_name(stream_name=stream_name)
Return the SQL table name of the named stream.
446 def get_sql_table( 447 self, 448 stream_name: str, 449 ) -> sqlalchemy.Table: 450 """Return a SQLAlchemy table object for the named stream.""" 451 return self.get_sql_cache().processor.get_sql_table(stream_name)
Return a SQLAlchemy table object for the named stream.
453 def get_dataset(self, stream_name: str) -> CachedDataset: 454 """Retrieve an `airbyte.datasets.CachedDataset` object for a given stream name. 455 456 This can be used to read and analyze the data in a SQL-based destination. 457 458 TODO: In a future iteration, we can consider providing stream configuration information 459 (catalog information) to the `CachedDataset` object via the "Get stream properties" 460 API: https://reference.airbyte.com/reference/getstreamproperties 461 """ 462 return CachedDataset( 463 self.get_sql_cache(), 464 stream_name=stream_name, 465 stream_configuration=False, # Don't look for stream configuration in cache. 466 )
Retrieve an airbyte.datasets.CachedDataset object for a given stream name.
This can be used to read and analyze the data in a SQL-based destination.
TODO: In a future iteration, we can consider providing stream configuration information
(catalog information) to the CachedDataset object via the "Get stream properties"
API: https://reference.airbyte.com/reference/getstreamproperties
468 def get_sql_database_name(self) -> str: 469 """Return the SQL database name.""" 470 cache = self.get_sql_cache() 471 return cache.get_database_name()
Return the SQL database name.
473 def get_sql_schema_name(self) -> str: 474 """Return the SQL schema name.""" 475 cache = self.get_sql_cache() 476 return cache.schema_name
Return the SQL schema name.
478 @property 479 def stream_names(self) -> list[str]: 480 """Return the set of stream names.""" 481 return self.connection.stream_names
Return the set of stream names.
483 @final 484 @property 485 def streams( 486 self, 487 ) -> _SyncResultStreams: # pyrefly: ignore[unknown-name] 488 """Return a mapping of stream names to `airbyte.CachedDataset` objects. 489 490 This is a convenience wrapper around the `stream_names` 491 property and `get_dataset()` method. 492 """ 493 return self._SyncResultStreams(self)
Return a mapping of stream names to airbyte.CachedDataset objects.
This is a convenience wrapper around the stream_names
property and get_dataset() method.
134@dataclass 135class SyncAttempt: 136 """Represents a single attempt of a sync job. 137 138 **This class is not meant to be instantiated directly.** Instead, obtain a `SyncAttempt` by 139 calling `.SyncResult.get_attempts()`. 140 """ 141 142 workspace: CloudWorkspace 143 connection: CloudConnection 144 job_id: int 145 attempt_number: int 146 _attempt_data: dict[str, Any] | None = None 147 148 @property 149 def attempt_id(self) -> int: 150 """Return the attempt ID.""" 151 return self._get_attempt_data()["id"] 152 153 @property 154 def status(self) -> str: 155 """Return the attempt status.""" 156 return self._get_attempt_data()["status"] 157 158 @property 159 def bytes_synced(self) -> int: 160 """Return the number of bytes synced in this attempt.""" 161 return self._get_attempt_data().get("bytesSynced", 0) 162 163 @property 164 def records_synced(self) -> int: 165 """Return the number of records synced in this attempt.""" 166 return self._get_attempt_data().get("recordsSynced", 0) 167 168 @property 169 def created_at(self) -> datetime: 170 """Return the creation time of the attempt.""" 171 timestamp = self._get_attempt_data()["createdAt"] 172 return ab_datetime_parse(timestamp) 173 174 def _get_attempt_data(self) -> dict[str, Any]: 175 """Get attempt data from the provided attempt data.""" 176 if self._attempt_data is None: 177 raise ValueError( 178 "Attempt data not provided. SyncAttempt should be created via " 179 "SyncResult.get_attempts()." 180 ) 181 return self._attempt_data["attempt"] 182 183 def get_full_log_text(self) -> str: 184 """Return the complete log text for this attempt. 185 186 Returns: 187 String containing all log text for this attempt, with lines separated by newlines. 188 """ 189 if self._attempt_data is None: 190 return "" 191 192 logs_data = self._attempt_data.get("logs") 193 if not logs_data: 194 return "" 195 196 result = "" 197 198 if "events" in logs_data: 199 log_events = logs_data["events"] 200 if log_events: 201 log_lines = [] 202 for event in log_events: 203 timestamp = event.get("timestamp", "") 204 level = event.get("level", "INFO") 205 message = event.get("message", "") 206 log_lines.append( 207 f"[{timestamp}] {level}: {message}" # pyrefly: ignore[bad-argument-type] 208 ) 209 result = "\n".join(log_lines) 210 elif "logLines" in logs_data: 211 log_lines = logs_data["logLines"] 212 if log_lines: 213 result = "\n".join(log_lines) 214 215 return result
Represents a single attempt of a sync job.
This class is not meant to be instantiated directly. Instead, obtain a SyncAttempt by
calling .SyncResult.get_attempts().
148 @property 149 def attempt_id(self) -> int: 150 """Return the attempt ID.""" 151 return self._get_attempt_data()["id"]
Return the attempt ID.
153 @property 154 def status(self) -> str: 155 """Return the attempt status.""" 156 return self._get_attempt_data()["status"]
Return the attempt status.
158 @property 159 def bytes_synced(self) -> int: 160 """Return the number of bytes synced in this attempt.""" 161 return self._get_attempt_data().get("bytesSynced", 0)
Return the number of bytes synced in this attempt.
163 @property 164 def records_synced(self) -> int: 165 """Return the number of records synced in this attempt.""" 166 return self._get_attempt_data().get("recordsSynced", 0)
Return the number of records synced in this attempt.
168 @property 169 def created_at(self) -> datetime: 170 """Return the creation time of the attempt.""" 171 timestamp = self._get_attempt_data()["createdAt"] 172 return ab_datetime_parse(timestamp)
Return the creation time of the attempt.
183 def get_full_log_text(self) -> str: 184 """Return the complete log text for this attempt. 185 186 Returns: 187 String containing all log text for this attempt, with lines separated by newlines. 188 """ 189 if self._attempt_data is None: 190 return "" 191 192 logs_data = self._attempt_data.get("logs") 193 if not logs_data: 194 return "" 195 196 result = "" 197 198 if "events" in logs_data: 199 log_events = logs_data["events"] 200 if log_events: 201 log_lines = [] 202 for event in log_events: 203 timestamp = event.get("timestamp", "") 204 level = event.get("level", "INFO") 205 message = event.get("message", "") 206 log_lines.append( 207 f"[{timestamp}] {level}: {message}" # pyrefly: ignore[bad-argument-type] 208 ) 209 result = "\n".join(log_lines) 210 elif "logLines" in logs_data: 211 log_lines = logs_data["logLines"] 212 if log_lines: 213 result = "\n".join(log_lines) 214 215 return result
Return the complete log text for this attempt.
Returns:
String containing all log text for this attempt, with lines separated by newlines.