Skip to content

feat(elt-common): SQL extract improvements - #406

Merged
WHTaylor merged 8 commits into
mainfrom
321-db-improvements
Jul 28, 2026
Merged

WHTaylor merged 8 commits into
mainfrom
321-db-improvements

Conversation

@WHTaylor

Copy link
Copy Markdown
Contributor

ref #321

A couple of improvments required for porting the opralogweb pipeline to elt-common.

  1. A configurable limit on the number of rows returned, mostly for testing purposes
  2. The ability to add arbitrary filters to a table query, to enable querying the MoreEntryColumns table
  3. Converting the DB table schema into a pyarrow schema. As explained in the commit comment, this is to avoid problems when all the values in a column are null (which is especially the case when testing with low row_limits).

@bashanlam can you try out the postgres pipeline with these changes incorporated to see if it works for that as well, or if there are more cases which need covering?

WHTaylor added 3 commits July 27, 2026 15:27
ref #321

This is necessary for cases where all the values in a given column (or
any single chunk of it) are null, because if that happens pyarrow
doesn't know what type the values are supposed to be. This either
crashes in the conversion to iceberg types, if we don't support that, or
breaks later when trying to upload something with an actual value,
because it looks like the type of the column has changed, which is
incompatible.
@WHTaylor
WHTaylor requested a review from a team as a code owner July 27, 2026 14:33
@coderabbitai

coderabbitai Bot commented Jul 27, 2026 •

Copy link
Copy Markdown
Contributor

Important

Review skipped

Auto incremental reviews are disabled on this repository.

Please check the settings in the CodeRabbit UI or the .coderabbit.yaml file in this repository. To trigger a single review, invoke the @coderabbitai review command.

⚙️ Run configuration

Configuration used: Path: .coderabbit.yaml

Review profile: CHILL

Plan: Pro Plus

Run ID: 9514eb5c-3af6-4b58-a554-614cf7b8f1ed

You can disable this status message by setting the reviews.review_status to false in the CodeRabbit configuration file.

Use the checkbox below for a quick retry:

  • 🔍 Trigger review
📝 Walkthrough

Walkthrough

Changes

SQL extraction enhancements

Layer / File(s) Summary
SQLAlchemy-to-PyArrow schema conversion
elt-common/src/elt_common/sources/sqldatabase/schema.py, elt-common/tests/unit_tests/sources/test_sqldatabase_schema.py
Adds SQLAlchemy type mappings, numeric handling, Python-type fallback, unsupported-type errors, and comprehensive schema tests.
Limited SQL extraction and explicit schemas
elt-common/src/elt_common/sources/sqldatabase/__init__.py, elt-common/tests/unit_tests/sources/test_sqldatabase.py
Adds configurable row limits and query filters, applies limits during extraction, and constructs Arrow partitions with explicit schemas tested across limit values.

Possibly related PRs

Suggested reviewers: martyngigg

🚥 Pre-merge checks | ✅ 5
✅ Passed checks (5 passed)
Check name Status Explanation
Title check ✅ Passed The title clearly summarises the main change: SQL extract improvements in elt-common.
Description check ✅ Passed The description is relevant and matches the row limits, query filters, and schema conversion changes.
Docstring Coverage ✅ Passed No functions found in the changed files to evaluate docstring coverage. Skipping docstring coverage check.
Linked Issues check ✅ Passed Check skipped because no linked issues were found for this pull request.
Out of Scope Changes check ✅ Passed Check skipped because no linked issues were found for this pull request.

Thanks for using CodeRabbit! It's free for OSS, and your support helps us grow. If you like it, consider giving us a shout-out.

❤️ Share

Comment @coderabbitai help to get the list of available commands.

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Actionable comments posted: 6

🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

Inline comments:
In `@elt-common/src/elt_common/sources/sqldatabase/__init__.py`:
- Line 156: Wire the configured query_filter from the resource-level TableInfo
(or equivalent per-resource configuration) through _make_table_properties() into
_extract_table() instead of always passing None. Preserve unfiltered extraction
when no filter is configured, and add an end-to-end extraction test covering a
configured SQL filter.
- Around line 35-36: Update the row_limit field configuration in the SQL
database source model to enforce a minimum value of 0 using a Pydantic v2 field
constraint or validator, while continuing to allow None for no limit. Add a
validation test covering negative row_limit values and preserve valid zero and
positive limits.

In `@elt-common/src/elt_common/sources/sqldatabase/schema.py`:
- Around line 85-90: Update the Numeric/NUMERIC handling in the SQL type
conversion function to explicitly handle precision above decimal128’s limit of
38: use pa.decimal256 when supported, otherwise raise a clear conversion error.
Preserve decimal128 for precision up to 38 and add coverage for Numeric(39,
scale).
- Around line 79-83: Update _to_pyarrow_type to resolve the direct _SQL_TYPE_MAP
entry before accessing sql_type.python_type, avoiding NotImplementedError for
unmapped or custom SQLAlchemy types. Only consult _PYTHON_TYPE_MAP through the
Python-type fallback when no direct mapping exists, and preserve the existing
factory invocation behavior.
- Line 55: Update the sa.JSON mapping in the SQL-to-Arrow schema conversion so
deserialized objects and arrays are normalized to text before conversion, or
mapped to an Arrow JSON-compatible schema instead of being treated directly as
strings. Preserve scalar JSON handling, and extend the JSON cases in
test_sqldatabase_schema.py to cover both a non-null object and array.
- Around line 24-25: Update
elt-common/src/elt_common/sources/sqldatabase/schema.py lines 24-25 to keep
sa.Double and sa.Float mapped to Arrow floating-point types, removing the
pa.decimal128 fallback. In _SQL_ROOT_TYPES at lines 49-51 and
_to_pyarrow_type/_SQL_TYPE_MAP handling at lines 72-73, ensure Decimal and float
do not invoke an unparameterized decimal128 factory before the type-map lookup;
preserve parameterized Decimal conversion.
🪄 Autofix (Beta)

Fix all unresolved CodeRabbit comments on this PR:

  • Push a commit to this branch (recommended)
  • Create a new PR with the fixes

ℹ️ Review info
⚙️ Run configuration

Configuration used: Path: .coderabbit.yaml

Review profile: CHILL

Plan: Pro Plus

Run ID: d1706059-ffe9-47e8-909f-f9a1917014a8

📥 Commits

Reviewing files that changed from the base of the PR and between e4057d8 and 35218d5.

📒 Files selected for processing (4)
  • elt-common/src/elt_common/sources/sqldatabase/__init__.py
  • elt-common/src/elt_common/sources/sqldatabase/schema.py
  • elt-common/tests/unit_tests/sources/test_sqldatabase.py
  • elt-common/tests/unit_tests/sources/test_sqldatabase_schema.py

Comment on lines +35 to +36
row_limit: Optional[int] = None
"""Maximum number of rows to return from each table, primarily for testing purposes. No limit if 'None'"""

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🎯 Functional Correctness | 🟠 Major | ⚡ Quick win

🧩 Analysis chain

🏁 Script executed:

#!/bin/bash
python - <<'PY'
import sqlite3

conn = sqlite3.connect(":memory:")
conn.execute("CREATE TABLE example (value INTEGER)")
conn.executemany("INSERT INTO example VALUES (?)", [(1,), (2,)])
print(conn.execute("SELECT value FROM example LIMIT -1").fetchall())
PY

Repository: ISISNeutronMuon/analytics-data-platform

Length of output: 190


🏁 Script executed:

#!/bin/bash
set -euo pipefail

# Locate and inspect the targeted SQLAlchemy database source model and related usages.
printf '--- files ---\n'
git ls-files | rg '(^|/)elt_common/sources/sqldatabase/__init__.py$|sqldatabase' | sed -n '1,120p'

printf '\n--- target file size ---\n'
wc -l elt-common/src/elt_common/sources/sqldatabase/__init__.py

printf '\n--- target file relevant sections ---\n'
sed -n '1,140p' elt-common/src/elt_common/sources/sqldatabase/__init__.py
printf '\n--- usages of SqlDatabaseConfig/row_limit ---\n'
rg -n "row_limit|SqlDatabaseConfig|Select\.limit|\.limit" elt-common/src -S

Repository: ISISNeutronMuon/analytics-data-platform

Length of output: 268


🏁 Script executed:

#!/bin/bash
set -euo pipefail

# Locate and inspect the targeted SQLAlchemy database source model and related usages.
printf '%s\n' '--- files ---'
git ls-files | rg '(^|/)elt_common/sources/sqldatabase/__init__.py$|sqldatabase' | sed -n '1,120p'

printf '%s\n' ''
printf '%s\n' '--- target file size ---'
wc -l elt-common/src/elt_common/sources/sqldatabase/__init__.py

printf '%s\n' ''
printf '%s\n' '--- target file relevant sections ---'
sed -n '1,180p' elt-common/src/elt_common/sources/sqldatabase/__init__.py

printf '%s\n' ''
printf '%s\n' '--- usages of SqlDatabaseConfig/row_limit ---'
rg -n "row_limit|SqlDatabaseConfig|Select\.limit|\.limit" elt-common/src -S

Repository: ISISNeutronMuon/analytics-data-platform

Length of output: 7394


🌐 Web query:

Python Pydantic v2 root_validator deprecated or model_validator use field and model_validator for class/instance validation

💡 Result:

In Pydantic V2, the @root_validator decorator is deprecated and has been replaced by the @model_validator decorator [1][2][3]. Use @model_validator for model-level validation, choosing the appropriate mode based on when the validation should occur: 1. mode='before': Runs before the model is instantiated [4][5]. The validator receives the raw input data (usually a dict) [4][5]. It must be defined as a @classmethod [4][5]. Example: from typing import Any from pydantic import BaseModel, model_validator class MyModel(BaseModel): data: str @model_validator(mode='before') @classmethod def validate_raw(cls, data: Any) -> Any: # Process raw data return data 2. mode='after': Runs after the model has been validated and instantiated [4][5]. The validator receives the model instance itself and should return the validated instance [4][5]. It is defined as an instance method (no @classmethod decorator) [4][5]. Example: from typing_extensions import Self from pydantic import BaseModel, model_validator class MyModel(BaseModel): value: int @model_validator(mode='after') def check_value(self) -> Self: # Access instance via self return self 3. mode='wrap': Wraps the validation process, allowing you to execute code before and after the inner validator runs [4][6]. Note that for field-specific validation, you should now use @field_validator instead of the legacy @validator decorator [1][2]. Migrating to these V2-style decorators ensures compatibility and unlocks new features, such as improved validation logic and better type safety [1][2].

Citations:


🏁 Script executed:

#!/bin/bash
set -euo pipefail

printf '%s\n' '--- project pydantic requirement ---'
rg -n "pydantic|pydantic-settings" pyproject.toml elk-common pylint.cfg requirements*.txt elt-common pyproject.toml elt-common/pyproject.toml -S 2>/dev/null || true

printf '%s\n' ''
printf '%s\n' '--- dependency files referencing pydantic ---'
for f in $(git ls-files | rg '(^|/)(pyproject\.toml|requirements.*\.txt|poetry\.lock|uv\.lock|Pipfile|Pipfile\.lock|conda.*\.ya?ml)$'); do
  if rg -q "pydantic|pydantic-settings" "$f"; then
    echo "--- $f ---"
    rg -n "pydantic|pydantic-settings" "$f" || true
  fi
done

printf '%s\n' ''
printf '%s\n' '--- tests around row_limit ---'
rg -n "row_limit|LIMIT|limit" elt-common/tests/unit_tests/sources/test_sqldatabase.py elt-common/tests/unit_tests/sources/test_sqldatabase_schema.py elt-common/tests -S || true

printf '%s\n' ''
printf '%s\n' '--- all in-repo row_limit assignments/usages ---'
rg -n "row_limit" -S .

Repository: ISISNeutronMuon/analytics-data-platform

Length of output: 50395


Reject negative SQL row limits during configuration validation.

row_limit is passed directly to SELECT.row_limit; on SQLite, negative values remove the LIMIT, so row_limit=-1 silently violates the documented maximum-row contract. Use a Pydantic v2 validator/field constraint such that row_limit >= 0 and include a negative-value test.

🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@elt-common/src/elt_common/sources/sqldatabase/__init__.py` around lines 35 -
36, Update the row_limit field configuration in the SQL database source model to
enforce a minimum value of 0 using a Pydantic v2 field constraint or validator,
while continuing to allow None for no limit. Add a validation test covering
negative row_limit values and preserve valid zero and positive limits.

Comment thread elt-common/src/elt_common/sources/sqldatabase/__init__.py
Comment on lines +24 to +25
sa.Double: pa.decimal128,
sa.Float: pa.float64,

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🎯 Functional Correctness | 🟠 Major | ⚡ Quick win

🧩 Analysis chain

🏁 Script executed:

#!/bin/bash
python - <<'PY'
import pyarrow as pa

try:
    pa.decimal128()
except TypeError as exc:
    print(type(exc).__name__, exc)
PY

Repository: ISISNeutronMuon/analytics-data-platform

Length of output: 297


🏁 Script executed:

#!/bin/bash
set -euo pipefail

echo "Schema file outline and relevant lines:"
if [ -f elt-common/src/elt_common/sources/sqldatabase/schema.py ]; then
  wc -l elt-common/src/elt_common/sources/sqldatabase/schema.py
  sed -n '1,120p' elt-common/src/elt_common/sources/sqldatabase/schema.py | cat -n
else
  echo "schema.py not found"
  fd -a 'schema.py$' .
fi

echo
echo "Search for _to_pyarrow_type usages:"
rg -n "_to_pyarrow_type|numeric128|decimal128|float64|sa\.Double|sa\.DECIMAL|sa\.DOUBLE" elt-common/src/elt_common/providers/schema.py elt-common/src/elt_common/sources/sqldatabase -g '*.py' || true

echo
echo "Inspect _to_pyarrow_type implementation around decimal branches:"
rg -n -C 4 "def _to_pyarrow_type|numeric128|decimal128|float64|Double" elt-common/src/elt_common/sources/sqldatabase -g '*.py' || true

Repository: ISISNeutronMuon/analytics-data-platform

Length of output: 9466


🌐 Web query:

PyArrow decimal128 signature requires precision Python API

💡 Result:

The pyarrow.decimal128 function in the PyArrow Python API has the following signature [1][2]: pyarrow.decimal128(int precision, int scale=0) Parameters: - precision (int): The number of significant digits the decimal type can represent. It must be between 1 and 38 [1][2]. - scale (int, default 0): The number of digits after the decimal point. This value can be negative [1][2]. This function creates a decimal type with a 128-bit width. Arrow decimals are represented as fixed-point numbers encoded as scaled integers [1][2]. If you require a precision greater than 38 significant digits, you should use pyarrow.decimal256 instead [1][2].

Citations:


Keep Double/FLOAT as Arrow floats and remove the unparameterisable pa.decimal128 fallback.

Decimal inherits _SQL_ROOT_TYPES[sa.Double], and float maps directly to pa.decimal128; both can end up calling factory() with no precision/scale before _to_pyarrow_type reaches the _SQL_TYPE_MAP lookup. The corrected fix is:

Proposed fix
-    sa.Double: pa.decimal128,
+    sa.Double: pa.float64,
...
-    sa.DECIMAL: _SQL_ROOT_TYPES[sa.Double],
     sa.DOUBLE: _SQL_ROOT_TYPES[sa.Double],
...
-    float: pa.decimal128,
+    float: pa.float64,
📍 Affects 1 file
  • elt-common/src/elt_common/sources/sqldatabase/schema.py#L24-L25 (this comment)
  • elt-common/src/elt_common/sources/sqldatabase/schema.py#L49-L51
  • elt-common/src/elt_common/sources/sqldatabase/schema.py#L72-L73
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@elt-common/src/elt_common/sources/sqldatabase/schema.py` around lines 24 -
25, Update elt-common/src/elt_common/sources/sqldatabase/schema.py lines 24-25
to keep sa.Double and sa.Float mapped to Arrow floating-point types, removing
the pa.decimal128 fallback. In _SQL_ROOT_TYPES at lines 49-51 and
_to_pyarrow_type/_SQL_TYPE_MAP handling at lines 72-73, ensure Decimal and float
do not invoke an unparameterized decimal128 factory before the type-map lookup;
preserve parameterized Decimal conversion.

sa.INTEGER: _SQL_ROOT_TYPES[sa.Integer],
sa.FLOAT: _SQL_ROOT_TYPES[sa.Float],
sa.INT: _SQL_ROOT_TYPES[sa.Integer],
sa.JSON: _SQL_ROOT_TYPES[sa.String],

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🗄️ Data Integrity & Integration | 🟠 Major | 🏗️ Heavy lift

🧩 Analysis chain

🏁 Script executed:

#!/bin/bash
python - <<'PY'
import pyarrow as pa

schema = pa.schema([pa.field("payload", pa.string())])
print(pa.Table.from_pylist([{"payload": {"key": "value"}}], schema=schema))
PY

Repository: ISISNeutronMuon/analytics-data-platform

Length of output: 297


🏁 Script executed:

#!/bin/bash
set -euo pipefail

# Locate and inspect the relevant schema file and nearby constants/usages.
file="elt-common/src/elt_common/sources/sqldatabase/schema.py"
if [ -f "$file" ]; then
  echo "== file exists =="
  wc -l "$file"
  echo "== relevant lines =="
  sed -n '1,140p' "$file" | nl -ba
else
  echo "missing $file"
  fd -a 'schema.py|pyproject.toml|requirements|setup.py' . | sed -n '1,80p'
fi

echo "== SQL JSON type mapping references =="
rg -n "sa\.JSON|_SQL_ROOT_TYPES|JSON" elt-common/src pyproject.toml requirements* setup.py 2>/dev/null || true

echo "== pyarrow availability =="
python3 - <<'PY'
import importlib.util
print("pyarrow_installed", bool(importlib.util.find_spec("pyarrow")))
PY

echo "== pyarrow package metadata if installed =="
python3 - <<'PY'
import importlib.util, json, subprocess
spec=importlib.util.find_spec("pyarrow")
print("pyarrow_loc", spec.origin if spec else None)
try:
    import pkg_resources
    print("pkg_resources_availability", True)
    req=pkg_resources.Requirement.parse("pyarrow")
    for dist in pkg_resources.working_set:
        if dist.project_name == "pyarrow":
            print(dist.as_requirement(), dist.version, dist.location)
            break
    else:
        print("pkg_resources_pyarrow_dist", None)
except Exception as e:
    print("pkg_resources_error", type(e).__name__, e)
PY

Repository: ISISNeutronMuon/analytics-data-platform

Length of output: 316


🏁 Script executed:

#!/bin/bash
set -euo pipefail

file="elt-common/src/elt_common/sources/sqldatabase/schema.py"
echo "== file exists =="
wc -l "$file"

echo "== relevant lines =="
awk '{printf "%6d\t%s\n", NR, $0}' "$file" | sed -n '1,120p'

echo "== SQL JSON type mapping references =="
rg -n "sa\.JSON|_SQL_ROOT_TYPES|JSON" elt-common/src pyproject.toml requirements* setup.py 2>/dev/null || true

echo "== pyarrow availability =="
python3 - <<'PY'
import importlib.util
print("pyarrow_installed", bool(importlib.util.find_spec("pyarrow")))
PY

Repository: ISISNeutronMuon/analytics-data-platform

Length of output: 7033


🏁 Script executed:

#!/bin/bash
set -euo pipefail

echo "== SQLAlchemy JSON python_type availability =="
python3 - <<'PY'
try:
    import sqlalchemy as sa
    json_type = sa.JSON()
    print("type:", type(json_type).__name__)
    print("python_type:", getattr(json_type, "python_type",))
    print("python_type type:", type(getattr(json_type, "python_type", None)).__name__)
    print("python_type is dict:", getattr(json_type, "python_type", None) is dict)
except Exception as e:
    print("error:", type(e).__name__, e)
PY

echo "== usages of to_pyarrow_schema / JSON columns in source =="
rg -n "to_pyarrow_schema|to_pyarrow_field|JSON|JSON columns|JSON to|arrow|convert" elt-common/src 2>/dev/null || true

echo "== tests mentioning JSON / schema conversion =="
fd -e py -d 3 . | grep -i test || true
rg -n "JSON|sqldatabase|to_pyarrow_schema|to_pyarrow_field" -g '*test*.py' -g '*.py' . 2>/dev/null | head -n 200 || true

Repository: ISISNeutronMuon/analytics-data-platform

Length of output: 10000


🌐 Web query:

SQLAlchemy JSON python_type dict list SQLAlchemy 2.0

💡 Result:

In SQLAlchemy 2.0, the JSON column type is typically mapped in Python as a dict or list, depending on the structure of the stored data [1]. Because standard Python dictionaries and lists are mutable, SQLAlchemy does not automatically detect in-place changes to these objects (e.g., modifying a key in a dictionary or appending to a list) [2][3]. To ensure these changes are tracked and saved during a session flush, you must use the sqlalchemy.ext.mutable extension [2][3]. 1. Using Mutation Tracking: To enable change tracking, you can use MutableDict or MutableList provided by the mutable extension [2]. For a dictionary-oriented structure: from sqlalchemy.ext.mutable import MutableDict from sqlalchemy.types import JSON from sqlalchemy.orm import mapped_column # Define the column using the mutable wrapper my_json_column = mapped_column(MutableDict.as_mutable(JSON)) 2. Type Annotations (Mapped): When using SQLAlchemy 2.0 Declarative mapping with Type Annotation Maps, you can define how dicts or lists should be treated by the ORM [1]: from typing import Any from sqlalchemy.orm import DeclarativeBase, Mapped from sqlalchemy.types import JSON class Base(DeclarativeBase): type_annotation_map = { dict[str, Any]: JSON, list[Any]: JSON, } class MyModel(Base): tablename = "my_table" # Now you can use Mapped[dict] or Mapped[list] directly data: Mapped[dict[str, Any]] = mapped_column 3. Important Considerations: * Nested Mutations: MutableDict and MutableList only track changes at the top level of the structure [2][4]. If your JSON data contains nested dictionaries or lists that you also need to track, standard MutableDict/MutableList will not suffice [2][3]. You would either need to implement a custom recursive mutable class or utilize third-party libraries like sqlalchemy-json to handle deep nested mutation tracking [3]. * python_type: In SQLAlchemy, the python_type property of a type refers to the expected Python type returned by the database driver for that column [5]. While historically this was often dict for JSON types, it is intended as a descriptor and does not inherently dictate ORM behavior [6][7]. For best practices in 2.0, rely on Declarative type annotation maps to define your ORM types [1].

Citations:


🌐 Web query:

PyArrow Table.from_pylist schema string non string value TypeError cast exception

💡 Result:

When using pyarrow.Table.from_pylist with an explicit schema, PyArrow strictly enforces that the input data matches the defined schema types [1][2]. A TypeError or ArrowTypeError occurring during this process typically indicates a failure to cast input values (such as a string) to the target schema type (such as an integer or timestamp) [3]. Key causes and solutions for this exception include: Incompatible Data Types: PyArrow cannot implicitly cast certain incompatible Python objects to target Arrow types [3]. For example, passing a string that does not represent a valid integer will cause a casting failure if the schema defines that column as pa.int64() [3]. Ensure your input data matches the schema's type requirements before conversion [3]. Incorrect Data Structure: from_pylist expects a list of dictionaries where each dictionary represents a row [1][4]. If you pass data structured for from_pydict (where keys are column names and values are lists of column data), the conversion will fail [3]. Verify that your input is a list of individual records, not a dictionary of columns [1][5]. Size and Memory Constraints: For extremely large inputs, PyArrow may implicitly convert data into ChunkedArray objects [6]. Certain operations or internal constructors may raise a TypeError if they specifically require a flat pyarrow.lib.Array and receive a ChunkedArray instead [7][6]. If you encounter this, verify if your data exceeds standard memory limits or array offset capacities (e.g., int32 limits) and consider processing the data in smaller batches or using pa.Table.from_pylist instead of pa.RecordBatch.from_pylist [7][6]. To debug, you can inspect the input data's first few rows to ensure the types match the expected schema [1]. You can also remove the schema argument temporarily to let PyArrow infer types, which will allow you to see what types PyArrow identifies automatically, helping you adjust your target schema accordingly [1][8].

Citations:


Do not declare deserialised JSON values as strings.

sa.JSON deserialises nested objects to dict/list, but this maps the column to pa.string(). Non-null JSON objects and arrays will fail the Arrow schema cast during table creation. Use a text-normalisation step before Arrow conversion, or an Arrow JSON-safe schema, and cover an object and an array in the test_sqldatabase_schema.py JSON tests.

🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@elt-common/src/elt_common/sources/sqldatabase/schema.py` at line 55, Update
the sa.JSON mapping in the SQL-to-Arrow schema conversion so deserialized
objects and arrays are normalized to text before conversion, or mapped to an
Arrow JSON-compatible schema instead of being treated directly as strings.
Preserve scalar JSON handling, and extend the JSON cases in
test_sqldatabase_schema.py to cover both a non-null object and array.

Comment thread elt-common/src/elt_common/sources/sqldatabase/schema.py
Comment on lines +85 to +90
if isinstance(sql_type, sa.Numeric) or isinstance(sql_type, sa.NUMERIC):
precision = getattr(sql_type, "precision", None)
scale = getattr(sql_type, "scale", None)
if precision is not None and scale is not None:
return pa.decimal128(precision, scale)
return pa.float64()

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🎯 Functional Correctness | 🟠 Major | ⚡ Quick win

🧩 Analysis chain

🏁 Script executed:

#!/bin/bash
python - <<'PY'
import pyarrow as pa

for precision in (38, 39):
    try:
        print(precision, pa.decimal128(precision, 2))
    except Exception as exc:
        print(precision, type(exc).__name__, exc)
PY

Repository: ISISNeutronMuon/analytics-data-platform

Length of output: 297


🏁 Script executed:

#!/bin/bash
set -euo pipefail

echo "== file and relevant lines =="
if [ -f elt-common/src/elt_common/sources/sqldatabase/schema.py ]; then
  sed -n '1,130p' elt-common/src/elt_common/sources/sqldatabase/schema.py
else
  echo "missing file"
fi

echo
echo "== pyarrow dependency/version declarations =="
rg -n "pyarrow|arrow" -S -g 'pyproject.toml' -g 'requirements*.txt' -g 'poetry.lock' -g 'Pipfile*' -g 'setup.py' -g 'setup.cfg' . || true

echo
echo "== searches for numeric/conversion behaviour =="
rg -n "sa\.Numeric|Numeric|decimal128|decimal256|scale|precision" elt-common/src/elt_common/sources/sqldatabase/schema.py elt-common/src -S || true

Repository: ISISNeutronMuon/analytics-data-platform

Length of output: 6899


🌐 Web query:

Apache Arrow Python decimal128 precision max 38 precision 39

💡 Result:

In Apache Arrow, the decimal128 type has a maximum supported precision of 38 significant digits [1][2]. This limit is enforced because 38 digits is the maximum number of decimal digits that can be safely and reliably represented within a 128-bit signed integer [3][4]. While a 128-bit signed integer technically has a range that can hold values requiring 39 digits (specifically, the values between 10^38 and 2^127 - 1), these values cannot be consistently represented or displayed as decimals across all operations without risk of overflow or loss of accuracy [5]. Consequently, the official Apache Arrow documentation and API implementations define the maximum precision as 38 [1][6][2]. If your application requires a precision greater than 38, you should use the decimal256 type, which supports up to 76 significant digits [1][6][2]. Note: In some older versions or specific implementations, the software may not have strictly validated the input precision, occasionally allowing values like 39 to be passed [7]. However, this is considered outside the defined specifications, and you should treat 38 as the hard limit for decimal128 to ensure compatibility and correctness [1][3][4].

Citations:


Support or reject decimal precision above 38 explicitly.

SQL Numeric(39, scale) reaches pa.decimal128(39, scale), which exceeds Arrow decimal128’s maximum precision. Select decimal256 where supported, or raise a clear conversion error, and cover this edge case.

🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@elt-common/src/elt_common/sources/sqldatabase/schema.py` around lines 85 -
90, Update the Numeric/NUMERIC handling in the SQL type conversion function to
explicitly handle precision above decimal128’s limit of 38: use pa.decimal256
when supported, otherwise raise a clear conversion error. Preserve decimal128
for precision up to 38 and add coverage for Numeric(39, scale).

WHTaylor added 5 commits July 27, 2026 15:55
CodeRabbit thinks this isn't as simple as coercing it to a string. It's simplest to leave it out for now, and handle it when we have a DB that has a JSON column
It may be possible for python_type to throw an error

@martyngigg martyngigg left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Nice addition of the row limit for testing.

It's a shame we have to manually deal with the schema but then again what is it supposed to do with no information! At least this also makes it clear that there are two schema conversions happening, SA -> PyArrow -> Iceberg, whereas before that was a bit hidden.

@bashanlam

Copy link
Copy Markdown
Contributor

@WHTaylor I tested the Postgres pipeline with these changes, but it failed because JSONB isn't supported in SQLAlchemy yet. If you can merge this to main, I can take care of the fix by updating schema.py (attached).
schema.py

@WHTaylor
WHTaylor merged commit e2558eb into main Jul 28, 2026
4 checks passed
@WHTaylor
WHTaylor deleted the 321-db-improvements branch July 28, 2026 09:41
bashanlam added a commit that referenced this pull request Jul 28, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants