feat(elt-common): Add Extractor creation - #348
Conversation
ref #321 Co-authored-by: Martyn Gigg <martyn.gigg@gmail.com>
|
Important Review skippedAuto incremental reviews are disabled on this repository. Please check the settings in the CodeRabbit UI or the ⚙️ Run configurationConfiguration used: Path: .coderabbit.yaml Review profile: CHILL Plan: Pro Run ID: You can disable this status message by setting the Use the checkbox below for a quick retry:
📝 WalkthroughWalkthroughThis PR introduces a factory-based extraction framework, allowing dynamic instantiation of extraction classes from job modules with automatic pydantic configuration binding. ChangesExtract factory framework and SQL integration
Possibly related PRs
Suggested reviewers
Poem
🚥 Pre-merge checks | ✅ 4 | ❌ 1❌ Failed checks (1 warning)
✅ Passed checks (4 passed)
✏️ 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. Comment |
|
@coderabbitai review |
✅ Action performedReview finished.
|
There was a problem hiding this comment.
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 winClosure captures
nameby reference, causing incorrect table extraction when extractors are called after iteration.The
extractorclosure captures the loop variablenameby reference (Python's late-binding behaviour). If a consumer collects all yielded properties before calling extractors, every extractor will use the final value ofname:# 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 valueEmpty 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 valueConsider prefixing unused parameter with underscore.
The
watermarkparameter is required by the extractor signature but is not used in the function body. Consider renaming it to_watermarkto 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 valueVariable name
iis misleading.Line 85 unpacks the tuple as
for i, props in ..., butireceives a resource name string (e.g., "0", "1", "2"), not an index. Consider renaming toresource_nameornamefor 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 valueConsider 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 valueConsider prefixing unused parameter with underscore.
The
watermarkparameter is required by the extractor signature but is not used in the function body. Consider renaming it to_watermarkto 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
📒 Files selected for processing (12)
elt-common/src/elt_common/extract.pyelt-common/src/elt_common/sources/sqldatabase/__init__.pyelt-common/tests/unit_tests/create_extract_obj_fakes/custom_config.pyelt-common/tests/unit_tests/create_extract_obj_fakes/errors/blank.pyelt-common/tests/unit_tests/create_extract_obj_fakes/errors/doesnt_implement_method.pyelt-common/tests/unit_tests/create_extract_obj_fakes/errors/non_class.pyelt-common/tests/unit_tests/create_extract_obj_fakes/errors/non_sub_class.pyelt-common/tests/unit_tests/create_extract_obj_fakes/errors/not_pythonelt-common/tests/unit_tests/create_extract_obj_fakes/sql_extract.pyelt-common/tests/unit_tests/create_extract_obj_fakes/three_empty_tables.pyelt-common/tests/unit_tests/sources/test_sqldatabase.pyelt-common/tests/unit_tests/test_extract.py
| 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" |
There was a problem hiding this comment.
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.
| 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 |
There was a problem hiding this comment.
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.
TIL python handles this automatically
martyngigg
left a comment
There was a problem hiding this comment.
Agreed offline that extract_resource_properties is a better indicator than resource_properties that something dynamic is happening.
There was a problem hiding this comment.
Thanks for picking up and testing the extract functionality thoroughly.
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>
ref #321
Adds functionality for creating
Extractorinstances from pipeline scripts. This is largely identical to how it exists onelt-command-without-dltwith some extra guardrails and testing.A couple of small structural changes:
resource_properties->extract_resource_propertiesas we discussedSqlDatabaseExtractextendBaseExtractso that it works with the new functionalityWith those in it should put
SqlDatabaseExtractinto basically its final form, so the class defined insql_extract.pyfor 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
Refactor