diff --git a/elt-common/src/elt_common/sources/sqldatabase/__init__.py b/elt-common/src/elt_common/sources/sqldatabase/__init__.py index 0773f276..8f9536eb 100644 --- a/elt-common/src/elt_common/sources/sqldatabase/__init__.py +++ b/elt-common/src/elt_common/sources/sqldatabase/__init__.py @@ -4,7 +4,7 @@ import logging from abc import abstractmethod from collections.abc import Callable, Generator, Iterable, Iterator -from typing import NamedTuple +from typing import Any, NamedTuple import pyarrow as pa import pyarrow.compute as pc @@ -191,6 +191,21 @@ def extractor(watermark, *, _name=name): yield resource_name, properties + def _columns_for_select(self, table: sa.Table) -> Iterable[Any]: + """Retrieve columns for selection applying any DB-specific requirements.""" + if self._engine.dialect.name == "oracle": + # Apply UTC extraction strictly to Oracle databases to avoid Thin mode DPY-3022 errors + selected_cols = [] + for col in table.columns: + col_type_str = str(col.type).upper() + if getattr(col.type, "timezone", False) or "WITH TIME ZONE" in col_type_str: + selected_cols.append(sa.func.sys_extract_utc(col).label(col.name)) + else: + selected_cols.append(col) + return selected_cols + else: + return table.columns + def _extract_table( self, name: str, @@ -205,7 +220,8 @@ def _extract_table( self._metadata, autoload_with=self._engine, ) - query = sa.select(table) + + query = sa.select(*self._columns_for_select(table)) if watermark is not None: column, max_value = watermark.column, watermark.value LOGGER.debug(f"Cursor value detected. Limiting query to {column} > {max_value}") @@ -214,11 +230,9 @@ def _extract_table( if query_filter: query = query_filter(query) - query = query.limit(self.config.row_limit) + if self.config.row_limit: + query = query.limit(self.config.row_limit) - # If all the values in a column are null pyarrow won't know what type - # the column should be, so we need to explicitly create a schema from - # the table pa_schema = to_pyarrow_schema(table) result = conn.execution_options(yield_per=self.config.chunk_size).execute(query) diff --git a/elt-common/src/elt_common/sources/sqldatabase/schema.py b/elt-common/src/elt_common/sources/sqldatabase/schema.py index 1b081969..aed9979a 100644 --- a/elt-common/src/elt_common/sources/sqldatabase/schema.py +++ b/elt-common/src/elt_common/sources/sqldatabase/schema.py @@ -6,6 +6,7 @@ import pyarrow as pa import sqlalchemy as sa from sqlalchemy.dialects import postgresql +from sqlalchemy.dialects.oracle import RAW def to_pyarrow_schema(table: sa.Table) -> pa.Schema: @@ -63,6 +64,7 @@ def _to_pyarrow_field(column: sa.Column) -> pa.Field: sa.VARCHAR: _SQL_ROOT_TYPES[sa.String], postgresql.JSON: _SQL_ROOT_TYPES[sa.JSON], postgresql.JSONB: _SQL_ROOT_TYPES[sa.JSON], + RAW: _SQL_ROOT_TYPES[sa.LargeBinary], } _SQL_TYPE_MAP = _SQL_ROOT_TYPES | _EXTENDED_SQL_TYPES diff --git a/elt-pipelines/fase/ingest/fase/isisuserdb/isisuserdb.py b/elt-pipelines/fase/ingest/fase/isisuserdb/isisuserdb.py new file mode 100644 index 00000000..99fd8cb3 --- /dev/null +++ b/elt-pipelines/fase/ingest/fase/isisuserdb/isisuserdb.py @@ -0,0 +1,36 @@ +import sqlalchemy as sa +from elt_common.extract import ResourceWriteProperties +from elt_common.sources.sqldatabase import ( + SqlDatabaseExtract, + SqlDatabaseSourceConfig, + TableInfo, +) + + +class PipelineOracleConfig(SqlDatabaseSourceConfig): + drivername: str = "oracle+oracledb" + tables: list[str] + + @property + def connection_url(self) -> sa.URL: + return sa.URL.create( + drivername=self.drivername, + username=self.username, + password=self.password.get_secret_value() if self.password else None, + host=self.host, + port=self.port, + query={"service_name": self.database}, + ) + + +class Extract(SqlDatabaseExtract): + config_cls = PipelineOracleConfig + + def table_info(self) -> dict[str, TableInfo | None]: + """Defines the target tables and their ingestion strategy.""" + return { + table_name: TableInfo( + write_properties=ResourceWriteProperties(write_mode="replace") + ) + for table_name in self.config.tables + } diff --git a/elt-pipelines/pyproject.toml b/elt-pipelines/pyproject.toml index 5bdc8562..ec61eae3 100644 --- a/elt-pipelines/pyproject.toml +++ b/elt-pipelines/pyproject.toml @@ -13,6 +13,11 @@ dependencies = [ proposal = [ "sqlalchemy[postgresql-psycopgbinary]>=2.0.0", ] + +isisuserdb = [ + "sqlalchemy[oracle-oracledb]>=2.0.0", +] + statusdisplay = [ "requests>=2.34.2", ] diff --git a/elt-pipelines/uv.lock b/elt-pipelines/uv.lock index 3e604567..dc991f8b 100644 --- a/elt-pipelines/uv.lock +++ b/elt-pipelines/uv.lock @@ -756,6 +756,9 @@ dependencies = [ ] [package.optional-dependencies] +isisuserdb = [ + { name = "sqlalchemy", extra = ["oracle-oracledb"] }, +] jira = [ { name = "atlassian-python-api" }, ] @@ -799,9 +802,10 @@ requires-dist = [ { name = "requests", marker = "extra == 'statusdisplay'", specifier = ">=2.34.2" }, { name = "scipy", marker = "extra == 'moderator-performance'", specifier = ">=1.17.1,<2" }, { name = "sqlalchemy", extras = ["mssql-pymssql"], marker = "extra == 'opralogweb'", specifier = ">=2.0.51" }, + { name = "sqlalchemy", extras = ["oracle-oracledb"], marker = "extra == 'isisuserdb'", specifier = ">=2.0.0" }, { name = "sqlalchemy", extras = ["postgresql-psycopgbinary"], marker = "extra == 'proposal'", specifier = ">=2.0.0" }, ] -provides-extras = ["proposal", "statusdisplay", "sharepoint", "opralogweb", "moderator-performance", "jira"] +provides-extras = ["proposal", "isisuserdb", "statusdisplay", "sharepoint", "opralogweb", "moderator-performance", "jira"] [package.metadata.requires-dev] dev = [{ name = "prek", specifier = ">=0.4.5" }] @@ -1634,6 +1638,42 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/ca/6f/a04e900f465ff3221ccc395522503e2d10e79fa21f2723c8e177aae1e0d1/opentelemetry_api-1.44.0-py3-none-any.whl", hash = "sha256:94b98c893a91b88657eaac1e3ba89618cdb85be6918196705354f34728b2cdef", size = 60018, upload-time = "2026-07-16T15:25:11.657Z" }, ] +[[package]] +name = "oracledb" +version = "26.0.0" +source = { registry = "https://pypi.org/simple" } +dependencies = [ + { name = "cryptography" }, + { name = "typing-extensions" }, +] +sdist = { url = "https://files.pythonhosted.org/packages/ad/0f/7eefd0e10328d94109180a877f9f1297b8b8dcc80f0bcccb0b0efea2ca9b/oracledb-26.0.0.tar.gz", hash = "sha256:d15728bae7546984203ae17f8296153d52d10ffdaf4929bebe46b70d5819d270", size = 912239, upload-time = "2026-09-10T16:39:01.602Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/52/09/0c388ff9bf418abad68f492780d3cef8a33ca8a109ccdc71ff0e06344319/oracledb-26.0.0-cp313-cp313-macosx_10_13_universal2.whl", hash = "sha256:f69bafe46bf854c70f19e2f1d068c1b9f3d3066ff79cf90354e0adfeafac1b08", size = 4907337, upload-time = "2026-09-10T16:39:46.657Z" }, + { url = "https://files.pythonhosted.org/packages/c5/74/f689dffb9abd3f008f081795677658ce7750371f4ec8718c78663648c19f/oracledb-26.0.0-cp313-cp313-manylinux2014_aarch64.manylinux_2_17_aarch64.manylinux_2_28_aarch64.whl", hash = "sha256:26b9934aa9a2c431e9e3126b524a8154b78f1f3a12c258a9e7ee311b78135a78", size = 2442648, upload-time = "2026-09-10T16:39:48.349Z" }, + { url = "https://files.pythonhosted.org/packages/fe/56/1bc5d3061ba4d35a70489757d2308865490b0a9abc5116cddfdf19954030/oracledb-26.0.0-cp313-cp313-manylinux2014_x86_64.manylinux_2_17_x86_64.manylinux_2_28_x86_64.whl", hash = "sha256:587acb63db873dc9784284e02b8c690730289255791bfa1e1ca70cbad50ac130", size = 2643188, upload-time = "2026-09-10T16:39:49.933Z" }, + { url = "https://files.pythonhosted.org/packages/44/09/db87288a174b1b03d36de305e96801f808e740d7986db96624ab5541c4eb/oracledb-26.0.0-cp313-cp313-musllinux_1_2_aarch64.whl", hash = "sha256:65275fc0f3e2c8f52549b6b8099d6eb5efe4748e6d62d4bd07bc3b006044aa60", size = 2497863, upload-time = "2026-09-10T16:39:51.392Z" }, + { url = "https://files.pythonhosted.org/packages/2a/78/82fd386fcc679eae6f65b2866ea96d8b06ae0eecf17b176f847e8ece9aa5/oracledb-26.0.0-cp313-cp313-musllinux_1_2_x86_64.whl", hash = "sha256:184869bbb43a1b074e4783b0da5b1793681e3129f4b02c305ad1ba4c6f931621", size = 2671060, upload-time = "2026-09-10T16:39:52.871Z" }, + { url = "https://files.pythonhosted.org/packages/ea/fe/f4c5825ccd7e8c221dad1d16405e02fb57cf7f965837ff74715fc865cbe6/oracledb-26.0.0-cp313-cp313-win32.whl", hash = "sha256:d82d5cea25884ab76db46173e1469588093f35e066a3e972abf1a6a91775edec", size = 1575954, upload-time = "2026-09-10T16:39:54.282Z" }, + { url = "https://files.pythonhosted.org/packages/3e/1a/7410ebf575438e2439827dc38be3851640b4406f880e69f78d92141892bc/oracledb-26.0.0-cp313-cp313-win_amd64.whl", hash = "sha256:747ec61dc933db33db4814b24c8cf0a57c7583482ca444b124e44c910f25bd54", size = 1969227, upload-time = "2026-09-10T16:39:55.578Z" }, + { url = "https://files.pythonhosted.org/packages/5f/dc/1a5bd6233d3de25be8efe48c96077078e7cf940154096dea789e17c8af88/oracledb-26.0.0-cp313-cp313-win_arm64.whl", hash = "sha256:85c775fd2b9f93c2a6ff48cabc2be307ba9fa951b7556e5149dc2ee62c1fbfbe", size = 1620768, upload-time = "2026-09-10T16:39:56.901Z" }, + { url = "https://files.pythonhosted.org/packages/1b/ad/72720d3a3066e0202deb4a4fb0640aa93c661ee0af486a23f4b8a5e92f4f/oracledb-26.0.0-cp314-cp314-macosx_10_15_universal2.whl", hash = "sha256:79b612ae836b2182d4c980e1670eb1b689cbe8a68fe5002104dc410fbe16e568", size = 4960030, upload-time = "2026-09-10T16:39:58.637Z" }, + { url = "https://files.pythonhosted.org/packages/e1/03/af2b24c0c099db043bc567d268adf194aa6db1f7867b9798932da5e8981d/oracledb-26.0.0-cp314-cp314-manylinux2014_aarch64.manylinux_2_17_aarch64.manylinux_2_28_aarch64.whl", hash = "sha256:3d24dfacbff11a4b3c15b72b446a56b8f83a4a4013ed694ab2719c79227dc81a", size = 2485508, upload-time = "2026-09-10T16:40:00.371Z" }, + { url = "https://files.pythonhosted.org/packages/25/30/3b913e044b04c72608cf00dc11c382ed8c3c11bf30114862adf0b0301e5a/oracledb-26.0.0-cp314-cp314-manylinux2014_x86_64.manylinux_2_17_x86_64.manylinux_2_28_x86_64.whl", hash = "sha256:4380a300f4094efe0282f9990929b9189a8e1aa52677d023dd9bdeb0163948e1", size = 2662863, upload-time = "2026-09-10T16:40:01.841Z" }, + { url = "https://files.pythonhosted.org/packages/b7/01/9d114dade656c6f4477ebdd1a1e3c3310ca675c6ae2af59b4fe466f58c33/oracledb-26.0.0-cp314-cp314-musllinux_1_2_aarch64.whl", hash = "sha256:fa365bb1f9016d1a688af06dd021481d63d49be00f4e5dc91cf7370d4a2adef8", size = 2542811, upload-time = "2026-09-10T16:40:03.495Z" }, + { url = "https://files.pythonhosted.org/packages/f4/9a/b634a2afc9dd025d93f25843a91a368a7726134b73999c55acc804b47f71/oracledb-26.0.0-cp314-cp314-musllinux_1_2_x86_64.whl", hash = "sha256:a4779752a3ea31558088a306d346964ef4107f8589266ce625651066f80972e7", size = 2690239, upload-time = "2026-09-10T16:40:05.509Z" }, + { url = "https://files.pythonhosted.org/packages/1d/80/cc0eeb56900792f6135d423b69133f2186869ccb9f998ab0d4eb2a5f2a8d/oracledb-26.0.0-cp314-cp314-win32.whl", hash = "sha256:ada309ac4a5f23f98d6b74627d1508f5fa47ad3a681fe42c7769d7ade281d04d", size = 1601816, upload-time = "2026-09-10T16:40:07.213Z" }, + { url = "https://files.pythonhosted.org/packages/d4/ac/c03f4fda9bda88fb694ccac78e6894331928d8b5c1244c34e8678cf644f5/oracledb-26.0.0-cp314-cp314-win_amd64.whl", hash = "sha256:7a92c2b002abf425078269c9d7db8544458c5dc6b59affb5a21ff075315d9468", size = 2026664, upload-time = "2026-09-10T16:40:08.684Z" }, + { url = "https://files.pythonhosted.org/packages/f6/01/399c6267811e38f0ac74c1a048ec1588431dcf80306392cf9da656f65f1c/oracledb-26.0.0-cp314-cp314-win_arm64.whl", hash = "sha256:3a77c5ab67ec3b2e2953bcca2b9fa31491b94c9985e02861e306e01ddc61737d", size = 1679554, upload-time = "2026-09-10T16:40:10.14Z" }, + { url = "https://files.pythonhosted.org/packages/cc/ff/44dd56f1f6f93f5d8efbfd6cf82235357d8408d36706664f46b8ccb57302/oracledb-26.0.0-cp315-cp315-macosx_10_15_universal2.whl", hash = "sha256:02b725959a714a911d0b03264d825ad9e631c68c9601e4d753b207e2245ce13d", size = 4958274, upload-time = "2026-09-10T16:40:11.701Z" }, + { url = "https://files.pythonhosted.org/packages/cb/b4/f1c3dae696a146b072332dd5bd8149e5d384a6089cf592dff2fe54aa6b50/oracledb-26.0.0-cp315-cp315-manylinux2014_aarch64.manylinux_2_17_aarch64.manylinux_2_28_aarch64.whl", hash = "sha256:98a38d5a82ae587f5d8dbe7c624eb59380fe5bb761c39ac7529cc6c5eafd83c8", size = 2487339, upload-time = "2026-09-10T16:40:13.422Z" }, + { url = "https://files.pythonhosted.org/packages/c1/cc/f202034fe4c8ecce3eac980b9f1c6bbd4e460bf660c2c248e22e56b84339/oracledb-26.0.0-cp315-cp315-manylinux2014_x86_64.manylinux_2_17_x86_64.manylinux_2_28_x86_64.whl", hash = "sha256:8ab05b254671ae1d12a3699f3cd9b9d28a6fba19db0e82842c38ea3bf587f9a1", size = 2683681, upload-time = "2026-09-10T16:40:15.318Z" }, + { url = "https://files.pythonhosted.org/packages/3c/bb/fe2035571ac695c80db6b5fcd8b2b18b0f9b8e20c913261b8f756876fc93/oracledb-26.0.0-cp315-cp315-musllinux_1_2_aarch64.whl", hash = "sha256:342a63514e59bc0860c0050bf9be56d0220ae3b0bd32a0e4d22a86dec532bc78", size = 2541907, upload-time = "2026-09-10T16:40:16.923Z" }, + { url = "https://files.pythonhosted.org/packages/51/de/6c8bfa247a9825b5511bca1e329468e3bb583f440d848b924e44eaedf594/oracledb-26.0.0-cp315-cp315-musllinux_1_2_x86_64.whl", hash = "sha256:262dacf9d90450e29f971ae68a8bcb81edad94194628c3f2d1afa50ec8b42140", size = 2712297, upload-time = "2026-09-10T16:40:18.336Z" }, + { url = "https://files.pythonhosted.org/packages/b2/b0/6c8af7bd9216ec55ba67a38ce6e10fd05fd399f884792d1b20e45efb3820/oracledb-26.0.0-cp315-cp315-win32.whl", hash = "sha256:64a9d42dbea41c1f0d45e890b57404e5fc1b1dfeb3c0939128ecaddac9eaacd6", size = 1600829, upload-time = "2026-09-10T16:40:19.969Z" }, + { url = "https://files.pythonhosted.org/packages/d9/7d/a3b51eefaaedbb49d23ea437d019d7eafcd2e65249be9090c1223543cb45/oracledb-26.0.0-cp315-cp315-win_amd64.whl", hash = "sha256:90794a9ba26ee6a6ccfaf47ae88b5d1bfe3b6192e3591fa6c0851f9ee9cabc7d", size = 2025725, upload-time = "2026-09-10T16:40:21.569Z" }, + { url = "https://files.pythonhosted.org/packages/f8/fb/48be38127faacd006c28ac2773cc1dbe41f42f89a38e79c8ec217bfba595/oracledb-26.0.0-cp315-cp315-win_arm64.whl", hash = "sha256:0fae924f6d07d33dbf267f52d686b0eeb79640853b749cdce1da513241eed1da", size = 1679377, upload-time = "2026-09-10T16:40:23.067Z" }, +] + [[package]] name = "orderly-set" version = "5.5.0" @@ -2756,6 +2796,9 @@ wheels = [ mssql-pymssql = [ { name = "pymssql" }, ] +oracle-oracledb = [ + { name = "oracledb" }, +] postgresql-psycopgbinary = [ { name = "psycopg", extra = ["binary"] }, ]