Fix MongoToS3Operator aggregate-pipeline detection before rendering - #70330
Merged
Conversation
mongo_query is a template field, rendered after __init__ runs. Caching is_pipeline = isinstance(self.mongo_query, list) in the constructor inspects the un-rendered value, so a query that resolves to a list only after templating would be sent through find() instead of aggregate(). Decide the aggregate-vs-find path in execute() from the rendered value instead.
MannXo
requested review from
amoghrajesh,
ashb,
bugraoz93,
gopidesupavan,
jason810496,
jscheffl,
o-nikolas and
potiuk
as code owners
July 23, 2026 21:01
78 tasks
shahar1
approved these changes
Jul 24, 2026
Contributor
Backport failed to create: v3-3-test. View the failure log Run detailsNote: As of Merging PRs targeted for Airflow 3.X In matter of doubt please ask in #release-management Slack channel.
You can attempt to backport this manually by running: cherry_picker c42540f v3-3-testThis should apply the commit to the v3-3-test branch and leave the commit in conflict state marking After you have resolved the conflicts, you can continue the backport process by running: cherry_picker --continueIf you don't have cherry-picker installed, see the installation guide. |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Compute the aggregate-vs-find decision in
execute()from the renderedmongo_query, instead of cachingself.is_pipeline = isinstance(self.mongo_query, list)in__init__.mongo_queryis a template field, so it is rendered after the constructor runs. Derivingis_pipelinein__init__inspects the un-rendered Jinja expression, so amongo_querythat resolves to a list only after templating would be misclassified and sent throughfind()instead ofaggregate(). Deciding inexecute()uses the rendered value.Part of the template-field constructor burn-down; removes the
MongoToS3Operatorentry fromscripts/ci/prek/validate_operators_init_exemptions.txtin the same PR, as the hook requires.related: #70296
Was generative AI tooling used to co-author this PR?
Generated-by: Claude Code (Opus 4.8) following the guidelines