Skip to content

feat(elt-common): Add Extractor creation - #348

Merged
WHTaylor merged 5 commits into
mainfrom
321-extract-classes
Jun 12, 2026
Merged

WHTaylor merged 5 commits into
mainfrom
321-extract-classes

Conversation

@WHTaylor

@WHTaylor WHTaylor commented Jun 12, 2026 •

Copy link
Copy Markdown
Contributor

ref #321

Adds functionality for creating Extractor instances from pipeline scripts. This is largely identical to how it exists on elt-command-without-dlt with some extra guardrails and testing.

A couple of small structural changes:

  • Renames resource_properties -> extract_resource_properties as we discussed
  • Makes SqlDatabaseExtract extend BaseExtract so that it works with the new functionality

With those in it should put SqlDatabaseExtract into basically its final form, so the class defined in sql_extract.py for testing should also serve as a decent skeleton for trying out the new structure with extracting FASE data @bashanlam.

Summary by CodeRabbit

Release Notes

  • New Features

    • Introduced a unified extraction framework for improved consistency across data sources
    • Added support for environment-based configuration of extraction jobs
  • Refactor

    • Updated SQL database extraction to align with the new extraction architecture
    • Enhanced dynamic loading capabilities for extraction definitions

@WHTaylor
WHTaylor requested a review from a team as a code owner June 12, 2026 12:10
@coderabbitai

coderabbitai Bot commented Jun 12, 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

Run ID: c36cb2b4-12c6-4224-952b-9ba09aa16424

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

This PR introduces a factory-based extraction framework, allowing dynamic instantiation of extraction classes from job modules with automatic pydantic configuration binding. SqlDatabaseExtract is refactored to integrate with this new BaseExtract base class, with method and attribute name standardisation. Comprehensive test fixtures and test coverage validate both error and success paths of the factory function.

Changes

Extract factory framework and SQL integration

Layer / File(s) Summary
BaseExtract abstract class and factory functions
elt-common/src/elt_common/extract.py
Introduces BaseExtract abstract base class with config instance storage and abstract extract_resource_properties() iterator contract. New factory function create_extract_obj() dynamically resolves an Extract class from a job's module file using importlib.util, validates it is a BaseExtract subclass, and instantiates it with pydantic-configured settings using a job-specific environment variable prefix. Helper functions _get_extract_cls(), _get_extract_cls_from_module_path(), and _import_module_from_path() encapsulate dynamic loading with explicit error handling for missing files, classes, wrong types, and invalid specs.
SqlDatabaseExtract refactored to BaseExtract
elt-common/src/elt_common/sources/sqldatabase/__init__.py
Changes SqlDatabaseExtract inheritance from ABC to BaseExtract, renames configuration attribute from source_config_cls to config_cls, updates constructor parameter from source_config to config, and renames the hook method from resource_properties() to extract_resource_properties(). Engine and metadata initialisation now use config values directly. Imports updated to pull BaseExtract from elt_common.extract.
Error-path test fixtures
elt-common/tests/unit_tests/create_extract_obj_fakes/errors/doesnt_implement_method.py, non_class.py, non_sub_class.py
Test fixture modules representing error scenarios: doesnt_implement_method.py provides an Extract class without implementing required abstract methods, non_class.py defines Extract as a string constant, and non_sub_class.py defines Extract as a plain class not inheriting BaseExtract.
Success-path test fixtures
elt-common/tests/unit_tests/create_extract_obj_fakes/custom_config.py, sql_extract.py, three_empty_tables.py
Test modules demonstrating correct BaseExtract implementations: custom_config.py defines a custom BaseSettings subclass ACustomConfig and an Extract yielding a single empty resource; sql_extract.py extends SqlDatabaseExtract with empty table_info(); three_empty_tables.py yields three resource entries, each with an empty extractor.
Factory and integration test suite
elt-common/tests/unit_tests/test_extract.py
Comprehensive test coverage of create_extract_obj() factory: parametrised tests verify exception types and messages for missing files, missing classes, non-class exports, and non-subclass types. Success-path tests confirm BaseExtract instance creation, environment-variable config binding for custom configurations, and SqlDatabaseExtract instantiation with chunk-size derivation from environment settings.
SQL extraction tests refactored
elt-common/tests/unit_tests/sources/test_sqldatabase.py
All existing SqlDatabaseExtract unit tests updated to call extract_resource_properties() instead of resource_properties(). Six test cases covering no-table-info scenarios, chunked reading, multiple tables, write-properties metadata, and watermark-filtered output remain functionally identical; only method name and iteration target updated.

Possibly related PRs

  • ISISNeutronMuon/analytics-data-platform#339: Directly refactors SqlDatabaseExtract from SQL ingestion implementation to subclass BaseExtract, switching wiring to config_cls, and renaming the generator method to match this framework.

Suggested reviewers

  • martyngigg

Poem

🐰 A factory sprouts from seeds of BaseExtract,
Dynamically loading extractors with perfect tact.
SQL tables now dance with the new abstract grace,
Config binds through the air with environment's embrace!

🚥 Pre-merge checks | ✅ 4 | ❌ 1

❌ Failed checks (1 warning)

Check name Status Explanation Resolution
Docstring Coverage ⚠️ Warning Docstring coverage is 20.69% which is insufficient. The required threshold is 80.00%. Write docstrings for the functions missing them to satisfy the coverage threshold.
✅ Passed checks (4 passed)
Check name Status Explanation
Description Check ✅ Passed Check skipped - CodeRabbit’s high-level summary is enabled.
Title check ✅ Passed The title clearly and concisely summarises the primary change: adding extractor creation functionality to the elt-common module.
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.

✏️ Tip: You can configure your own custom pre-merge checks in the settings.


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 and usage tips.

@WHTaylor

Copy link
Copy Markdown
Contributor Author

@coderabbitai review

@coderabbitai

coderabbitai Bot commented Jun 12, 2026 •

Copy link
Copy Markdown
Contributor
✅ Action performed

Review finished.

Note: CodeRabbit is an incremental review system and does not re-review already reviewed commits. This command is applicable only when automatic reviews are paused.

@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: 3

Caution

Some comments are outside the diff and can’t be posted inline due to platform limitations.

⚠️ Outside diff range comments (1)
elt-common/src/elt_common/sources/sqldatabase/__init__.py (1)

132-133: ⚠️ Potential issue | 🟠 Major | ⚡ Quick win

Closure captures name by reference, causing incorrect table extraction when extractors are called after iteration.

The extractor closure captures the loop variable name by reference (Python's late-binding behaviour). If a consumer collects all yielded properties before calling extractors, every extractor will use the final value of name:

# This usage pattern will extract the wrong tables:
all_props = list(e.extract_resource_properties())
for table_name, props in all_props:
    data = props.extractor(None)  # All use last table name!

Current tests call extractors immediately within the iteration loop, so they don't expose this bug.

Proposed fix using default argument capture
-            def extractor(watermark):
-                return self._extract_table(name, watermark=watermark, conn=conn)
+            def extractor(watermark, *, _name=name):
+                return self._extract_table(_name, watermark=watermark, conn=conn)
🤖 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 132 -
133, The extractor closure in extract_resource_properties captures the loop
variable name by reference, causing all extractors to use the final name when
called later; fix it by binding name as a default parameter in the closure
(e.g., change def extractor(watermark): ... to def extractor(watermark,
name=name): return self._extract_table(name, watermark=watermark, conn=conn)),
ensuring each extractor calls _extract_table with the correct table name.
🧹 Nitpick comments (5)
elt-common/tests/unit_tests/create_extract_obj_fakes/custom_config.py (2)

17-25: 💤 Low value

Empty string as resource name may cause confusion.

Line 19 uses an empty string "" as the resource name. Whilst this may be valid, it could make debugging more difficult. Consider using a descriptive name like "test_resource" for clarity, especially since other fixtures (e.g., three_empty_tables.py) use meaningful identifiers.

🤖 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/tests/unit_tests/create_extract_obj_fakes/custom_config.py` around
lines 17 - 25, The resource name in extract_resource_properties currently yields
an empty string which is confusing; update the tuple yielded in
extract_resource_properties to use a descriptive identifier (e.g.,
"test_resource") instead of "" so test output and debugging are clearer—leave
the ResourceProperties construction (extractor=extract_nothing,
write_properties=ResourceWriteProperties(), watermark_column=None) untouched.

28-29: 💤 Low value

Consider prefixing unused parameter with underscore.

The watermark parameter is required by the extractor signature but is not used in the function body. Consider renaming it to _watermark to signal that it's intentionally unused.

♻️ Proposed change
-def extract_nothing(watermark):
+def extract_nothing(_watermark):
     yield []
🤖 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/tests/unit_tests/create_extract_obj_fakes/custom_config.py` around
lines 28 - 29, The extractor function extract_nothing defines an unused
parameter watermark; rename that parameter to _watermark in the extract_nothing
signature to indicate it's intentionally unused, update any internal references
(none expected) and any tests or call sites that reference the parameter name if
they destructure or rely on the exact signature, and run tests to ensure the
signature change doesn't break introspection or tooling that expects the
original name.
elt-common/tests/unit_tests/test_extract.py (2)

79-88: 💤 Low value

Variable name i is misleading.

Line 85 unpacks the tuple as for i, props in ..., but i receives a resource name string (e.g., "0", "1", "2"), not an index. Consider renaming to resource_name or name for clarity.

♻️ Proposed change
-    for i, props in extract_obj.extract_resource_properties():
+    for resource_name, props in extract_obj.extract_resource_properties():
         yielded = [d for d in props.extractor(None)]
         assert len(yielded) == 1
         assert yielded[0] == []
🤖 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/tests/unit_tests/test_extract.py` around lines 79 - 88, Rename the
misleading loop variable in test_create_extract_obj: when iterating over
create_extract_obj(job).extract_resource_properties(), the first element is a
resource name string, not an index, so change the variable from i to a clearer
name like resource_name (or name) to improve readability and avoid confusion
with numeric indices; update the loop header in test_create_extract_obj and any
references inside the loop accordingly, keeping the call to
extract_resource_properties() and the assertions unchanged.

111-111: 💤 Low value

Consider testing public interface instead of private attribute.

Line 111 directly accesses the private attribute _chunk_size. Whilst this is acceptable for integration testing, it couples the test to implementation details. If a public method or property becomes available for accessing chunk size, prefer that instead.

🤖 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/tests/unit_tests/test_extract.py` at line 111, The test currently
inspects the private attribute extract_obj._chunk_size; instead assert via the
public API—either use a public property like extract_obj.chunk_size if
available, or exercise the chunking behavior by calling a public method such as
extract_obj.chunks(...) or extract_obj.get_chunk_size() and asserting the
resulting chunk sizes equal 100; update the test to reference
extract_obj.chunk_size or validate output from extract_obj.chunks/get_chunk_size
rather than _chunk_size.
elt-common/tests/unit_tests/create_extract_obj_fakes/three_empty_tables.py (1)

19-20: 💤 Low value

Consider prefixing unused parameter with underscore.

The watermark parameter is required by the extractor signature but is not used in the function body. Consider renaming it to _watermark to signal that it's intentionally unused.

♻️ Proposed change
-def extract_empty(watermark):
+def extract_empty(_watermark):
     yield []
🤖 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/tests/unit_tests/create_extract_obj_fakes/three_empty_tables.py`
around lines 19 - 20, Rename the unused parameter in the extractor function
extract_empty from watermark to _watermark to signal it is intentionally unused;
update the function signature def extract_empty(_watermark): and leave the body
yielding [] unchanged so the extractor still matches the required signature
while avoiding linter warnings about an unused variable.
🤖 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/extract.py`:
- Around line 128-135: The exception construction for missing extract class is
wrong: instead of passing the caught exception `e` as a second argument to
AttributeError, re-raise a new AttributeError using exception chaining (raise
AttributeError(f"Module '{module_name}' doesn't include an '{EXTRACT_CLS_NAME}'
class, which is required for defining an ingest job") from e) so the original
traceback is preserved; update the code around the getattr(module,
EXTRACT_CLS_NAME) error handling (the try/except that sets extract_cls) to use
the "raise ... from e" form.

In `@elt-common/tests/unit_tests/test_extract.py`:
- Around line 101-111: The test test_create_extract_obj_sql_extract sets
environment variables (SQL_EXTRACT__DRIVERNAME, SQL_EXTRACT__DATABASE,
SQL_EXTRACT__CHUNK_SIZE) but doesn't clean them up, risking test pollution;
update the test to use pytest's monkeypatch fixture (or explicitly save and
restore os.environ) to set these keys only for the test's duration before
calling create_extract_obj, ensuring the environment is restored afterward so
subsequent tests using BaseExtract/SqlDatabaseExtract are unaffected.
- Around line 91-98: The test test_create_extract_obj_custom_config leaves
CUSTOM_CONFIG__REQUIRED_STR in the environment; modify the test to ensure
cleanup by using pytest's monkeypatch fixture (call
monkeypatch.setenv("CUSTOM_CONFIG__REQUIRED_STR", "required") before calling
create_extract_obj) or wrap the env set in try/finally and del
os.environ["CUSTOM_CONFIG__REQUIRED_STR"] in the finally block so the
environment is restored after the assertion involving create_extract_obj and
BaseExtract.

---

Outside diff comments:
In `@elt-common/src/elt_common/sources/sqldatabase/__init__.py`:
- Around line 132-133: The extractor closure in extract_resource_properties
captures the loop variable name by reference, causing all extractors to use the
final name when called later; fix it by binding name as a default parameter in
the closure (e.g., change def extractor(watermark): ... to def
extractor(watermark, name=name): return self._extract_table(name,
watermark=watermark, conn=conn)), ensuring each extractor calls _extract_table
with the correct table name.

---

Nitpick comments:
In `@elt-common/tests/unit_tests/create_extract_obj_fakes/custom_config.py`:
- Around line 17-25: The resource name in extract_resource_properties currently
yields an empty string which is confusing; update the tuple yielded in
extract_resource_properties to use a descriptive identifier (e.g.,
"test_resource") instead of "" so test output and debugging are clearer—leave
the ResourceProperties construction (extractor=extract_nothing,
write_properties=ResourceWriteProperties(), watermark_column=None) untouched.
- Around line 28-29: The extractor function extract_nothing defines an unused
parameter watermark; rename that parameter to _watermark in the extract_nothing
signature to indicate it's intentionally unused, update any internal references
(none expected) and any tests or call sites that reference the parameter name if
they destructure or rely on the exact signature, and run tests to ensure the
signature change doesn't break introspection or tooling that expects the
original name.

In `@elt-common/tests/unit_tests/create_extract_obj_fakes/three_empty_tables.py`:
- Around line 19-20: Rename the unused parameter in the extractor function
extract_empty from watermark to _watermark to signal it is intentionally unused;
update the function signature def extract_empty(_watermark): and leave the body
yielding [] unchanged so the extractor still matches the required signature
while avoiding linter warnings about an unused variable.

In `@elt-common/tests/unit_tests/test_extract.py`:
- Around line 79-88: Rename the misleading loop variable in
test_create_extract_obj: when iterating over
create_extract_obj(job).extract_resource_properties(), the first element is a
resource name string, not an index, so change the variable from i to a clearer
name like resource_name (or name) to improve readability and avoid confusion
with numeric indices; update the loop header in test_create_extract_obj and any
references inside the loop accordingly, keeping the call to
extract_resource_properties() and the assertions unchanged.
- Line 111: The test currently inspects the private attribute
extract_obj._chunk_size; instead assert via the public API—either use a public
property like extract_obj.chunk_size if available, or exercise the chunking
behavior by calling a public method such as extract_obj.chunks(...) or
extract_obj.get_chunk_size() and asserting the resulting chunk sizes equal 100;
update the test to reference extract_obj.chunk_size or validate output from
extract_obj.chunks/get_chunk_size rather than _chunk_size.
🪄 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

Run ID: 32990547-a6f9-4f4c-ba92-3e96d52e30b2

📥 Commits

Reviewing files that changed from the base of the PR and between c500d2c and ce927ee.

📒 Files selected for processing (12)
  • elt-common/src/elt_common/extract.py
  • elt-common/src/elt_common/sources/sqldatabase/__init__.py
  • elt-common/tests/unit_tests/create_extract_obj_fakes/custom_config.py
  • elt-common/tests/unit_tests/create_extract_obj_fakes/errors/blank.py
  • elt-common/tests/unit_tests/create_extract_obj_fakes/errors/doesnt_implement_method.py
  • elt-common/tests/unit_tests/create_extract_obj_fakes/errors/non_class.py
  • elt-common/tests/unit_tests/create_extract_obj_fakes/errors/non_sub_class.py
  • elt-common/tests/unit_tests/create_extract_obj_fakes/errors/not_python
  • elt-common/tests/unit_tests/create_extract_obj_fakes/sql_extract.py
  • elt-common/tests/unit_tests/create_extract_obj_fakes/three_empty_tables.py
  • elt-common/tests/unit_tests/sources/test_sqldatabase.py
  • elt-common/tests/unit_tests/test_extract.py

Comment thread elt-common/src/elt_common/extract.py
Comment on lines +91 to +98
def test_create_extract_obj_custom_config():
job = make_manifest("custom_config")
os.environ["CUSTOM_CONFIG__REQUIRED_STR"] = "required"

extract_obj = create_extract_obj(job)

assert isinstance(extract_obj, BaseExtract)
assert getattr(extract_obj.config, "required_str") == "required"

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.

⚠️ Potential issue | 🟡 Minor | ⚡ Quick win

Environment variable not cleaned up after test.

Line 93 sets os.environ["CUSTOM_CONFIG__REQUIRED_STR"] but doesn't remove it after the test completes. This could cause test pollution if tests run in a shared environment.

🧹 Proposed fix using pytest's monkeypatch fixture
-def test_create_extract_obj_custom_config():
+def test_create_extract_obj_custom_config(monkeypatch):
     job = make_manifest("custom_config")
-    os.environ["CUSTOM_CONFIG__REQUIRED_STR"] = "required"
+    monkeypatch.setenv("CUSTOM_CONFIG__REQUIRED_STR", "required")

     extract_obj = create_extract_obj(job)

     assert isinstance(extract_obj, BaseExtract)
     assert getattr(extract_obj.config, "required_str") == "required"

Alternatively, use a try-finally block or pytest fixture to ensure cleanup.

🤖 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/tests/unit_tests/test_extract.py` around lines 91 - 98, The test
test_create_extract_obj_custom_config leaves CUSTOM_CONFIG__REQUIRED_STR in the
environment; modify the test to ensure cleanup by using pytest's monkeypatch
fixture (call monkeypatch.setenv("CUSTOM_CONFIG__REQUIRED_STR", "required")
before calling create_extract_obj) or wrap the env set in try/finally and del
os.environ["CUSTOM_CONFIG__REQUIRED_STR"] in the finally block so the
environment is restored after the assertion involving create_extract_obj and
BaseExtract.

Comment on lines +101 to +111
def test_create_extract_obj_sql_extract():
job = make_manifest("sql_extract")
os.environ["SQL_EXTRACT__DRIVERNAME"] = "sqlite"
os.environ["SQL_EXTRACT__DATABASE"] = "not_real"
os.environ["SQL_EXTRACT__CHUNK_SIZE"] = "100"

extract_obj = create_extract_obj(job)

assert isinstance(extract_obj, BaseExtract)
assert isinstance(extract_obj, SqlDatabaseExtract)
assert extract_obj._chunk_size == 100

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.

⚠️ Potential issue | 🟡 Minor | ⚡ Quick win

Environment variables not cleaned up after test.

Lines 103-105 set multiple environment variables but don't remove them after the test completes. This could cause test pollution if tests run in a shared environment.

🧹 Proposed fix using pytest's monkeypatch fixture
-def test_create_extract_obj_sql_extract():
+def test_create_extract_obj_sql_extract(monkeypatch):
     job = make_manifest("sql_extract")
-    os.environ["SQL_EXTRACT__DRIVERNAME"] = "sqlite"
-    os.environ["SQL_EXTRACT__DATABASE"] = "not_real"
-    os.environ["SQL_EXTRACT__CHUNK_SIZE"] = "100"
+    monkeypatch.setenv("SQL_EXTRACT__DRIVERNAME", "sqlite")
+    monkeypatch.setenv("SQL_EXTRACT__DATABASE", "not_real")
+    monkeypatch.setenv("SQL_EXTRACT__CHUNK_SIZE", "100")

     extract_obj = create_extract_obj(job)

     assert isinstance(extract_obj, BaseExtract)
     assert isinstance(extract_obj, SqlDatabaseExtract)
     assert extract_obj._chunk_size == 100
🤖 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/tests/unit_tests/test_extract.py` around lines 101 - 111, The test
test_create_extract_obj_sql_extract sets environment variables
(SQL_EXTRACT__DRIVERNAME, SQL_EXTRACT__DATABASE, SQL_EXTRACT__CHUNK_SIZE) but
doesn't clean them up, risking test pollution; update the test to use pytest's
monkeypatch fixture (or explicitly save and restore os.environ) to set these
keys only for the test's duration before calling create_extract_obj, ensuring
the environment is restored afterward so subsequent tests using
BaseExtract/SqlDatabaseExtract are unaffected.

@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.

Agreed offline that extract_resource_properties is a better indicator than resource_properties that something dynamic is happening.

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.

Thanks for picking up and testing the extract functionality thoroughly.

@WHTaylor
WHTaylor merged commit 31b800e into main Jun 12, 2026
4 checks passed
@WHTaylor
WHTaylor deleted the 321-extract-classes branch June 12, 2026 15:44
WHTaylor added a commit that referenced this pull request Jun 18, 2026
ref #321 

Adds the `runner` module which orchestrates the other parts of the
package.

Largely the same as on `elt-command-without-dlt`, with some changes to
fit the structural changes mentioned in #348

I've also linked it up with the CLI to test running it against a 'real'
file, and it was able to load 'data' into my local iceberg catalog.
Interacting with iceberg is still _very_ slow, unfortunately - might
have to see if there are any pyiceberg settings to tweak to try and help
with that.

<!-- This is an auto-generated comment: release notes by coderabbit.ai
-->
## Summary by CodeRabbit

## Release Notes

* **New Features**
* Ingestion runner can now extract data from job-defined extractors and
write it to Iceberg, including per-table merge/partition/sort handling.
* Watermark support is added, with automatic protection against
out-of-order batches.
* CLI now accepts both qualified and unqualified job names, and defaults
`step` to `"all"`.

* **Bug Fixes**
* Replace-mode ingestion now correctly follows with append for
subsequent chunks.

* **Tests**
* Added end-to-end coverage for write modes, watermark behaviour, and
property generation.
<!-- end of auto-generated comment: release notes by coderabbit.ai -->

---------

Co-authored-by: Martyn Gigg <martyn.gigg@gmail.com>
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.

2 participants