Rewrite decorated task mapping - #21328

Merged
uranusjr merged 6 commits into
apache:mainfrom
astronomer:decorator-unmap
Feb 10, 2022
Merged

Rewrite decorated task mapping#21328
uranusjr merged 6 commits into
apache:mainfrom
astronomer:decorator-unmap

Conversation

@uranusjr

@uranusjruranusjr commented Feb 4, 2022

Copy link
Copy Markdown
Member

Need to get #21210 merged first and rebase this. Done

This rewrites how map() is implemented on _TaskDecorator. The previous implementation incorrectly assumed mapping a task decorator works like mapping a traditional operator, but in fact they have entirely different semantics.

When mapping a traditional operator, FooOperator.map(my_arg=[1, 2, 3]), the argument on FooOperator is mapped. But when mapping a task decorator, my_task.map(my_val=[1, 2, 3]), it’s the argument on the function wrapped by the task object being mapped. Therefore, mapped (and also partial-ed) values for my_val should go into the DecoratedOperator’s op_kwargs (and op_args) arguments instead.

An end-to-end test is provided to validate the simplest case of this. Also fixed a test where the end-to-end test becomes broken after running a DAG (because I need to run a second DAG!)

@boring-cyborgboring-cyborgBot added area:CLI provider:cncf-kubernetes Kubernetes (k8s) provider related issues area:Scheduler including HA (high availability) scheduler area:serialization labels Feb 4, 2022
@uranusjr
uranusjr marked this pull request as ready for review February 4, 2022 14:32
@uranusjr
uranusjrforce-pushed the decorator-unmap branch 6 times, most recently from b52938c to f963a48CompareFebruary 6, 2022 12:11
@ashb

ashb commented Feb 7, 2022

Copy link
Copy Markdown
Member

The serialization needs some work for this approach. Here is an incomplete test:

deftest_mapped_decorator_serde():
fromairflow.models.xcom_argimportXComArgfromairflow.decoratorsimporttaskwithDAG("test-dag", start_date=datetime(2020, 1, 1)) asdag:
task1=BaseOperator(task_id="op1")
xcomarg=XComArg(task1, "test_key")
@task(retry_delay=30)defx(arg1, arg2):
...
real_op=x.partial(arg1=1).map(arg2=xcomarg).operatorserialized=SerializedBaseOperator._serialize(real_op)
assertserialized== {
'_is_dummy': False,
'_is_mapped': True,
'_task_module': 'airflow.decorators.python',
'_task_type': '_PythonDecoratedOperator',
'downstream_task_ids': [],
'partial_kwargs': {
'multiple_outputs': False,
'op_args': [],
'op_kwargs': {
'arg1': [
1,
2,
{"__type": "dict", "__var": {'a': 'b'}},
],
},
'retry_delay': 30,
},
'mapped_kwargs': {
'op_args': [],
# We don't need the __type/__var here!'op_kwargs': {
'arg2': {'__type': 'xcomref', '__var': {'task_id': 'op1', 'key': 'test_key'}},
}
},
'task_id': 'x',
'template_ext': [],
'template_fields': ['op_args', 'op_kwargs'],
# We don't want to include the python source code in the serialized representation# TODO? Where does `retry_delay` go? We need to separate 
}
op=SerializedBaseOperator.deserialize_operator(serialized)
assertisinstance(op, MappedOperator)
assertop.depsisMappedOperator.DEFAULT_DEPS# TODO: add some more asserts here

The serialization

Trying to call real_op.unmap() also throws a key error:

tests/serialization/test_dag_serialization.py:1670: in test_mapped_decorator_serde
real_op.unmap()
airflow/models/baseoperator.py:1882: in unmap
dag._remove_task(self.task_id)
airflow/models/dag.py:2167: in _remove_task
task = self.task_dict.pop(task_id)
E KeyError: 'x'

@uranusjr

Copy link
Copy Markdown
MemberAuthor

Serialisation format is revised. The KeyError is also fixed by implementing slightly smarter logic to merge op_kwargs from mapped_kwargs and partial_kwargs.

@uranusjr

Copy link
Copy Markdown
MemberAuthor

We should also implement some validation when a DAG is parsed to make sure the user pass reasonable values to partial_kwargs. This will be done in a separate PR.

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.

Suggested change
assertdeserialized.retry_delay==timedelta(seconds=30)

(give or take.)

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

At this point retry_delay is still 30 verbatim; it become a timedelta only after unmapped. I’ll add a test elsewhere for this.

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.

It should be deserialized as a timedelta -- it's fine for this value since the constructor for BaseOperator handles it, but other variables might behave differently

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.

start_date for instance.

@uranusjruranusjrFeb 9, 2022

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

I added a test in tests/decorators/test_python.py for this

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

start_date for instance.

Hm some refactoring would be called for to extract the logic out of BaseOperator for reuse. (Note that this affects non-decorator MappedOperator as well.) I think this should be done in a separate PR.

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.

Yeah agreed. We can fix this up later, this unblocks a lot.

Mapping a traditional operator (where arguments go to the operator)
and a task flow operator (where arguments go to the *function*) have
very different semantics, so we need some special code for them.
Previously we were merging partial_kwargs and mapped_kwargs too naively
and did not correctly handle op_args and op_kwargs; those need special
logic due to the mapping semantics of decorated tasks.
Some attributes are removed from serialization to match the format
of the (unmapped) _PythonDecoratedOperator. Some simplication is
implemented to op_kwargs to save some space.
ashb
ashb approved these changes Feb 9, 2022
@github-actions

Copy link
Copy Markdown
Contributor

The PR most likely needs to run full matrix of tests because it modifies parts of the core of Airflow. However, committers might decide to merge it quickly and take the risk. If they don't merge it quickly - please rebase it to the latest main at your convenience, or amend the last commit of the PR, and push it with --force-with-lease.

@github-actionsgithub-actionsBot added the full tests needed We need to run full set of tests for this PR to merge label Feb 9, 2022
@uranusjr

Copy link
Copy Markdown
MemberAuthor

Static check failures fixed in #21480.

@uranusjr
uranusjr merged commit fded2ca into apache:mainFeb 10, 2022
@uranusjr
uranusjr deleted the decorator-unmap branch February 10, 2022 07:07
ferruzzi pushed a commit to ferruzzi/airflow that referenced this pull request Feb 11, 2022
@jedcunninghamjedcunningham added changelog:skip Changes that should be skipped from the changelog (CI, tests, etc..) area:dynamic-task-mapping AIP-42 labels Feb 28, 2022
@jedcunninghamjedcunningham added this to the Airflow 2.3.0 milestone Apr 26, 2022
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:CLIarea:dynamic-task-mappingAIP-42area:Schedulerincluding HA (high availability) schedulerarea:serializationchangelog:skipChanges that should be skipped from the changelog (CI, tests, etc..)full tests neededWe need to run full set of tests for this PR to mergeprovider:cncf-kubernetesKubernetes (k8s) provider related issues

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants

@uranusjr@ashb@jedcunningham
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Add copy buttons to all
 blocks\n(function() {\n function addCopyButtons() {\n document.querySelectorAll('pre code').forEach(function(codeBlock) {\n if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;\n codeBlock.parentElement.setAttribute('data-copy-added', 'true');\n \n var btn = document.createElement('button');\n btn.textContent = 'Copy';\n btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';\n btn.onmouseover = function() { this.style.opacity = '1'; };\n btn.onmouseout = function() { this.style.opacity = '0.7'; };\n btn.onclick = function() {\n navigator.clipboard.writeText(codeBlock.textContent).then(function() {\n btn.textContent = 'Copied!';\n setTimeout(function() { btn.textContent = 'Copy'; }, 1500);\n });\n };\n codeBlock.parentElement.style.position = 'relative';\n codeBlock.parentElement.appendChild(btn);\n });\n }\n \n addCopyButtons();\n \n // Re-run on dynamic content\n var observer = new MutationObserver(addCopyButtons);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Add Copy Buttons to Code Blocks");
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
Skip to content

Rewrite decorated task mapping - #21328

Merged
uranusjr merged 6 commits into
apache:mainfrom
astronomer:decorator-unmap
Feb 10, 2022
Merged

Rewrite decorated task mapping#21328
uranusjr merged 6 commits into
apache:mainfrom
astronomer:decorator-unmap

Conversation

@uranusjr

@uranusjruranusjr commented Feb 4, 2022

Copy link
Copy Markdown
Member

Need to get #21210 merged first and rebase this. Done

This rewrites how map() is implemented on _TaskDecorator. The previous implementation incorrectly assumed mapping a task decorator works like mapping a traditional operator, but in fact they have entirely different semantics.

When mapping a traditional operator, FooOperator.map(my_arg=[1, 2, 3]), the argument on FooOperator is mapped. But when mapping a task decorator, my_task.map(my_val=[1, 2, 3]), it’s the argument on the function wrapped by the task object being mapped. Therefore, mapped (and also partial-ed) values for my_val should go into the DecoratedOperator’s op_kwargs (and op_args) arguments instead.

An end-to-end test is provided to validate the simplest case of this. Also fixed a test where the end-to-end test becomes broken after running a DAG (because I need to run a second DAG!)

@boring-cyborgboring-cyborgBot added area:CLI provider:cncf-kubernetes Kubernetes (k8s) provider related issues area:Scheduler including HA (high availability) scheduler area:serialization labels Feb 4, 2022
@uranusjr
uranusjr marked this pull request as ready for review February 4, 2022 14:32
@uranusjr
uranusjrforce-pushed the decorator-unmap branch 6 times, most recently from b52938c to f963a48CompareFebruary 6, 2022 12:11
@ashb

ashb commented Feb 7, 2022

Copy link
Copy Markdown
Member

The serialization needs some work for this approach. Here is an incomplete test:

deftest_mapped_decorator_serde():
fromairflow.models.xcom_argimportXComArgfromairflow.decoratorsimporttaskwithDAG("test-dag", start_date=datetime(2020, 1, 1)) asdag:
task1=BaseOperator(task_id="op1")
xcomarg=XComArg(task1, "test_key")
@task(retry_delay=30)defx(arg1, arg2):
...
real_op=x.partial(arg1=1).map(arg2=xcomarg).operatorserialized=SerializedBaseOperator._serialize(real_op)
assertserialized== {
'_is_dummy': False,
'_is_mapped': True,
'_task_module': 'airflow.decorators.python',
'_task_type': '_PythonDecoratedOperator',
'downstream_task_ids': [],
'partial_kwargs': {
'multiple_outputs': False,
'op_args': [],
'op_kwargs': {
'arg1': [
1,
2,
{"__type": "dict", "__var": {'a': 'b'}},
],
},
'retry_delay': 30,
},
'mapped_kwargs': {
'op_args': [],
# We don't need the __type/__var here!'op_kwargs': {
'arg2': {'__type': 'xcomref', '__var': {'task_id': 'op1', 'key': 'test_key'}},
}
},
'task_id': 'x',
'template_ext': [],
'template_fields': ['op_args', 'op_kwargs'],
# We don't want to include the python source code in the serialized representation# TODO? Where does `retry_delay` go? We need to separate 
}
op=SerializedBaseOperator.deserialize_operator(serialized)
assertisinstance(op, MappedOperator)
assertop.depsisMappedOperator.DEFAULT_DEPS# TODO: add some more asserts here

The serialization

Trying to call real_op.unmap() also throws a key error:

tests/serialization/test_dag_serialization.py:1670: in test_mapped_decorator_serde
real_op.unmap()
airflow/models/baseoperator.py:1882: in unmap
dag._remove_task(self.task_id)
airflow/models/dag.py:2167: in _remove_task
task = self.task_dict.pop(task_id)
E KeyError: 'x'

@uranusjr

Copy link
Copy Markdown
MemberAuthor

Serialisation format is revised. The KeyError is also fixed by implementing slightly smarter logic to merge op_kwargs from mapped_kwargs and partial_kwargs.

@uranusjr

Copy link
Copy Markdown
MemberAuthor

We should also implement some validation when a DAG is parsed to make sure the user pass reasonable values to partial_kwargs. This will be done in a separate PR.

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.

Suggested change
assertdeserialized.retry_delay==timedelta(seconds=30)

(give or take.)

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

At this point retry_delay is still 30 verbatim; it become a timedelta only after unmapped. I’ll add a test elsewhere for this.

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.

It should be deserialized as a timedelta -- it's fine for this value since the constructor for BaseOperator handles it, but other variables might behave differently

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.

start_date for instance.

@uranusjruranusjrFeb 9, 2022

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

I added a test in tests/decorators/test_python.py for this

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

start_date for instance.

Hm some refactoring would be called for to extract the logic out of BaseOperator for reuse. (Note that this affects non-decorator MappedOperator as well.) I think this should be done in a separate PR.

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.

Yeah agreed. We can fix this up later, this unblocks a lot.

Mapping a traditional operator (where arguments go to the operator)
and a task flow operator (where arguments go to the *function*) have
very different semantics, so we need some special code for them.
Previously we were merging partial_kwargs and mapped_kwargs too naively
and did not correctly handle op_args and op_kwargs; those need special
logic due to the mapping semantics of decorated tasks.
Some attributes are removed from serialization to match the format
of the (unmapped) _PythonDecoratedOperator. Some simplication is
implemented to op_kwargs to save some space.
ashb
ashb approved these changes Feb 9, 2022
@github-actions

Copy link
Copy Markdown
Contributor

The PR most likely needs to run full matrix of tests because it modifies parts of the core of Airflow. However, committers might decide to merge it quickly and take the risk. If they don't merge it quickly - please rebase it to the latest main at your convenience, or amend the last commit of the PR, and push it with --force-with-lease.

@github-actionsgithub-actionsBot added the full tests needed We need to run full set of tests for this PR to merge label Feb 9, 2022
@uranusjr

Copy link
Copy Markdown
MemberAuthor

Static check failures fixed in #21480.

@uranusjr
uranusjr merged commit fded2ca into apache:mainFeb 10, 2022
@uranusjr
uranusjr deleted the decorator-unmap branch February 10, 2022 07:07
ferruzzi pushed a commit to ferruzzi/airflow that referenced this pull request Feb 11, 2022
@jedcunninghamjedcunningham added changelog:skip Changes that should be skipped from the changelog (CI, tests, etc..) area:dynamic-task-mapping AIP-42 labels Feb 28, 2022
@jedcunninghamjedcunningham added this to the Airflow 2.3.0 milestone Apr 26, 2022
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:CLIarea:dynamic-task-mappingAIP-42area:Schedulerincluding HA (high availability) schedulerarea:serializationchangelog:skipChanges that should be skipped from the changelog (CI, tests, etc..)full tests neededWe need to run full set of tests for this PR to mergeprovider:cncf-kubernetesKubernetes (k8s) provider related issues

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants

@uranusjr@ashb@jedcunningham
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Force GitHub README to respect dark mode\n(function() {\n var style = document.createElement('style');\n style.textContent = '\n .markdown-body {\n color-scheme: dark light;\n }\n .markdown-body pre { background: #161b22 !important; }\n .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; }\n .markdown-body table th, .markdown-body table td { border-color: #30363d !important; }\n .markdown-body img { background: #0d1117; }\n .markdown-body blockquote { border-left-color: #8b949e; }\n .markdown-body hr { border-color: #30363d; }\n ';\n document.head.appendChild(style);\n})();", "GitHub Dark Mode README Fix"); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Rewrite decorated task mapping - #21328

Merged
uranusjr merged 6 commits into
apache:mainfrom
astronomer:decorator-unmap
Feb 10, 2022
Merged

Rewrite decorated task mapping#21328
uranusjr merged 6 commits into
apache:mainfrom
astronomer:decorator-unmap

Conversation

@uranusjr

@uranusjruranusjr commented Feb 4, 2022

Copy link
Copy Markdown
Member

Need to get #21210 merged first and rebase this. Done

This rewrites how map() is implemented on _TaskDecorator. The previous implementation incorrectly assumed mapping a task decorator works like mapping a traditional operator, but in fact they have entirely different semantics.

When mapping a traditional operator, FooOperator.map(my_arg=[1, 2, 3]), the argument on FooOperator is mapped. But when mapping a task decorator, my_task.map(my_val=[1, 2, 3]), it’s the argument on the function wrapped by the task object being mapped. Therefore, mapped (and also partial-ed) values for my_val should go into the DecoratedOperator’s op_kwargs (and op_args) arguments instead.

An end-to-end test is provided to validate the simplest case of this. Also fixed a test where the end-to-end test becomes broken after running a DAG (because I need to run a second DAG!)

@boring-cyborgboring-cyborgBot added area:CLI provider:cncf-kubernetes Kubernetes (k8s) provider related issues area:Scheduler including HA (high availability) scheduler area:serialization labels Feb 4, 2022
@uranusjr
uranusjr marked this pull request as ready for review February 4, 2022 14:32
@uranusjr
uranusjrforce-pushed the decorator-unmap branch 6 times, most recently from b52938c to f963a48CompareFebruary 6, 2022 12:11
@ashb

ashb commented Feb 7, 2022

Copy link
Copy Markdown
Member

The serialization needs some work for this approach. Here is an incomplete test:

deftest_mapped_decorator_serde():
fromairflow.models.xcom_argimportXComArgfromairflow.decoratorsimporttaskwithDAG("test-dag", start_date=datetime(2020, 1, 1)) asdag:
task1=BaseOperator(task_id="op1")
xcomarg=XComArg(task1, "test_key")
@task(retry_delay=30)defx(arg1, arg2):
...
real_op=x.partial(arg1=1).map(arg2=xcomarg).operatorserialized=SerializedBaseOperator._serialize(real_op)
assertserialized== {
'_is_dummy': False,
'_is_mapped': True,
'_task_module': 'airflow.decorators.python',
'_task_type': '_PythonDecoratedOperator',
'downstream_task_ids': [],
'partial_kwargs': {
'multiple_outputs': False,
'op_args': [],
'op_kwargs': {
'arg1': [
1,
2,
{"__type": "dict", "__var": {'a': 'b'}},
],
},
'retry_delay': 30,
},
'mapped_kwargs': {
'op_args': [],
# We don't need the __type/__var here!'op_kwargs': {
'arg2': {'__type': 'xcomref', '__var': {'task_id': 'op1', 'key': 'test_key'}},
}
},
'task_id': 'x',
'template_ext': [],
'template_fields': ['op_args', 'op_kwargs'],
# We don't want to include the python source code in the serialized representation# TODO? Where does `retry_delay` go? We need to separate 
}
op=SerializedBaseOperator.deserialize_operator(serialized)
assertisinstance(op, MappedOperator)
assertop.depsisMappedOperator.DEFAULT_DEPS# TODO: add some more asserts here

The serialization

Trying to call real_op.unmap() also throws a key error:

tests/serialization/test_dag_serialization.py:1670: in test_mapped_decorator_serde
real_op.unmap()
airflow/models/baseoperator.py:1882: in unmap
dag._remove_task(self.task_id)
airflow/models/dag.py:2167: in _remove_task
task = self.task_dict.pop(task_id)
E KeyError: 'x'

@uranusjr

Copy link
Copy Markdown
MemberAuthor

Serialisation format is revised. The KeyError is also fixed by implementing slightly smarter logic to merge op_kwargs from mapped_kwargs and partial_kwargs.

@uranusjr

Copy link
Copy Markdown
MemberAuthor

We should also implement some validation when a DAG is parsed to make sure the user pass reasonable values to partial_kwargs. This will be done in a separate PR.

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.

Suggested change
assertdeserialized.retry_delay==timedelta(seconds=30)

(give or take.)

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

At this point retry_delay is still 30 verbatim; it become a timedelta only after unmapped. I’ll add a test elsewhere for this.

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.

It should be deserialized as a timedelta -- it's fine for this value since the constructor for BaseOperator handles it, but other variables might behave differently

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.

start_date for instance.

@uranusjruranusjrFeb 9, 2022

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

I added a test in tests/decorators/test_python.py for this

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

start_date for instance.

Hm some refactoring would be called for to extract the logic out of BaseOperator for reuse. (Note that this affects non-decorator MappedOperator as well.) I think this should be done in a separate PR.

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.

Yeah agreed. We can fix this up later, this unblocks a lot.

Mapping a traditional operator (where arguments go to the operator)
and a task flow operator (where arguments go to the *function*) have
very different semantics, so we need some special code for them.
Previously we were merging partial_kwargs and mapped_kwargs too naively
and did not correctly handle op_args and op_kwargs; those need special
logic due to the mapping semantics of decorated tasks.
Some attributes are removed from serialization to match the format
of the (unmapped) _PythonDecoratedOperator. Some simplication is
implemented to op_kwargs to save some space.
ashb
ashb approved these changes Feb 9, 2022
@github-actions

Copy link
Copy Markdown
Contributor

The PR most likely needs to run full matrix of tests because it modifies parts of the core of Airflow. However, committers might decide to merge it quickly and take the risk. If they don't merge it quickly - please rebase it to the latest main at your convenience, or amend the last commit of the PR, and push it with --force-with-lease.

@github-actionsgithub-actionsBot added the full tests needed We need to run full set of tests for this PR to merge label Feb 9, 2022
@uranusjr

Copy link
Copy Markdown
MemberAuthor

Static check failures fixed in #21480.

@uranusjr
uranusjr merged commit fded2ca into apache:mainFeb 10, 2022
@uranusjr
uranusjr deleted the decorator-unmap branch February 10, 2022 07:07
ferruzzi pushed a commit to ferruzzi/airflow that referenced this pull request Feb 11, 2022
@jedcunninghamjedcunningham added changelog:skip Changes that should be skipped from the changelog (CI, tests, etc..) area:dynamic-task-mapping AIP-42 labels Feb 28, 2022
@jedcunninghamjedcunningham added this to the Airflow 2.3.0 milestone Apr 26, 2022
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:CLIarea:dynamic-task-mappingAIP-42area:Schedulerincluding HA (high availability) schedulerarea:serializationchangelog:skipChanges that should be skipped from the changelog (CI, tests, etc..)full tests neededWe need to run full set of tests for this PR to mergeprovider:cncf-kubernetesKubernetes (k8s) provider related issues

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants

@uranusjr@ashb@jedcunningham
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Highlight search terms from Google/DuckDuckGo/Bing referrer\n(function() {\n var ref = document.referrer;\n var terms = [];\n \n if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) {\n var url = new URL(ref);\n var q = url.searchParams.get('q') || url.searchParams.get('p');\n if (q) {\n terms = q.split(/\\s+/).filter(function(t) { return t.length > 2; });\n }\n }\n \n if (terms.length === 0) return;\n \n var style = document.createElement('style');\n style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }';\n document.head.appendChild(style);\n \n function highlight(node) {\n if (node.nodeType === 3) { // text node\n var text = node.textContent;\n var found = false;\n terms.forEach(function(term) {\n var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\\]\\\\]/g, '\\\\') + ')', 'gi');\n if (regex.test(text)) {\n found = true;\n var frag = document.createDocumentFragment();\n var parts = text.split(regex);\n parts.forEach(function(part, i) {\n if (i % 2 === 0) {\n frag.appendChild(document.createTextNode(part));\n } else {\n var span = document.createElement('span');\n span.className = 'userscript-highlight';\n span.textContent = part;\n frag.appendChild(span);\n }\n });\n node.parentNode.replaceChild(frag, node);\n }\n });\n } else if (node.nodeType === 1 && node.childNodes) { // element\n var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT'];\n if (!skipTags.includes(node.tagName)) {\n Array.from(node.childNodes).forEach(highlight);\n }\n }\n }\n \n highlight(document.body);\n \n // Re-highlight on dynamic content\n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1 || node.nodeType === 3) highlight(node);\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Highlight Search Terms"); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Rewrite decorated task mapping - #21328

Merged
uranusjr merged 6 commits into
apache:mainfrom
astronomer:decorator-unmap
Feb 10, 2022
Merged

Rewrite decorated task mapping#21328
uranusjr merged 6 commits into
apache:mainfrom
astronomer:decorator-unmap

Conversation

@uranusjr

@uranusjruranusjr commented Feb 4, 2022

Copy link
Copy Markdown
Member

Need to get #21210 merged first and rebase this. Done

This rewrites how map() is implemented on _TaskDecorator. The previous implementation incorrectly assumed mapping a task decorator works like mapping a traditional operator, but in fact they have entirely different semantics.

When mapping a traditional operator, FooOperator.map(my_arg=[1, 2, 3]), the argument on FooOperator is mapped. But when mapping a task decorator, my_task.map(my_val=[1, 2, 3]), it’s the argument on the function wrapped by the task object being mapped. Therefore, mapped (and also partial-ed) values for my_val should go into the DecoratedOperator’s op_kwargs (and op_args) arguments instead.

An end-to-end test is provided to validate the simplest case of this. Also fixed a test where the end-to-end test becomes broken after running a DAG (because I need to run a second DAG!)

@boring-cyborgboring-cyborgBot added area:CLI provider:cncf-kubernetes Kubernetes (k8s) provider related issues area:Scheduler including HA (high availability) scheduler area:serialization labels Feb 4, 2022
@uranusjr
uranusjr marked this pull request as ready for review February 4, 2022 14:32
@uranusjr
uranusjrforce-pushed the decorator-unmap branch 6 times, most recently from b52938c to f963a48CompareFebruary 6, 2022 12:11
@ashb

ashb commented Feb 7, 2022

Copy link
Copy Markdown
Member

The serialization needs some work for this approach. Here is an incomplete test:

deftest_mapped_decorator_serde():
fromairflow.models.xcom_argimportXComArgfromairflow.decoratorsimporttaskwithDAG("test-dag", start_date=datetime(2020, 1, 1)) asdag:
task1=BaseOperator(task_id="op1")
xcomarg=XComArg(task1, "test_key")
@task(retry_delay=30)defx(arg1, arg2):
...
real_op=x.partial(arg1=1).map(arg2=xcomarg).operatorserialized=SerializedBaseOperator._serialize(real_op)
assertserialized== {
'_is_dummy': False,
'_is_mapped': True,
'_task_module': 'airflow.decorators.python',
'_task_type': '_PythonDecoratedOperator',
'downstream_task_ids': [],
'partial_kwargs': {
'multiple_outputs': False,
'op_args': [],
'op_kwargs': {
'arg1': [
1,
2,
{"__type": "dict", "__var": {'a': 'b'}},
],
},
'retry_delay': 30,
},
'mapped_kwargs': {
'op_args': [],
# We don't need the __type/__var here!'op_kwargs': {
'arg2': {'__type': 'xcomref', '__var': {'task_id': 'op1', 'key': 'test_key'}},
}
},
'task_id': 'x',
'template_ext': [],
'template_fields': ['op_args', 'op_kwargs'],
# We don't want to include the python source code in the serialized representation# TODO? Where does `retry_delay` go? We need to separate 
}
op=SerializedBaseOperator.deserialize_operator(serialized)
assertisinstance(op, MappedOperator)
assertop.depsisMappedOperator.DEFAULT_DEPS# TODO: add some more asserts here

The serialization

Trying to call real_op.unmap() also throws a key error:

tests/serialization/test_dag_serialization.py:1670: in test_mapped_decorator_serde
real_op.unmap()
airflow/models/baseoperator.py:1882: in unmap
dag._remove_task(self.task_id)
airflow/models/dag.py:2167: in _remove_task
task = self.task_dict.pop(task_id)
E KeyError: 'x'

@uranusjr

Copy link
Copy Markdown
MemberAuthor

Serialisation format is revised. The KeyError is also fixed by implementing slightly smarter logic to merge op_kwargs from mapped_kwargs and partial_kwargs.

@uranusjr

Copy link
Copy Markdown
MemberAuthor

We should also implement some validation when a DAG is parsed to make sure the user pass reasonable values to partial_kwargs. This will be done in a separate PR.

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.

Suggested change
assertdeserialized.retry_delay==timedelta(seconds=30)

(give or take.)

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

At this point retry_delay is still 30 verbatim; it become a timedelta only after unmapped. I’ll add a test elsewhere for this.

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.

It should be deserialized as a timedelta -- it's fine for this value since the constructor for BaseOperator handles it, but other variables might behave differently

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.

start_date for instance.

@uranusjruranusjrFeb 9, 2022

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

I added a test in tests/decorators/test_python.py for this

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

start_date for instance.

Hm some refactoring would be called for to extract the logic out of BaseOperator for reuse. (Note that this affects non-decorator MappedOperator as well.) I think this should be done in a separate PR.

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.

Yeah agreed. We can fix this up later, this unblocks a lot.

Mapping a traditional operator (where arguments go to the operator)
and a task flow operator (where arguments go to the *function*) have
very different semantics, so we need some special code for them.
Previously we were merging partial_kwargs and mapped_kwargs too naively
and did not correctly handle op_args and op_kwargs; those need special
logic due to the mapping semantics of decorated tasks.
Some attributes are removed from serialization to match the format
of the (unmapped) _PythonDecoratedOperator. Some simplication is
implemented to op_kwargs to save some space.
ashb
ashb approved these changes Feb 9, 2022
@github-actions

Copy link
Copy Markdown
Contributor

The PR most likely needs to run full matrix of tests because it modifies parts of the core of Airflow. However, committers might decide to merge it quickly and take the risk. If they don't merge it quickly - please rebase it to the latest main at your convenience, or amend the last commit of the PR, and push it with --force-with-lease.

@github-actionsgithub-actionsBot added the full tests needed We need to run full set of tests for this PR to merge label Feb 9, 2022
@uranusjr

Copy link
Copy Markdown
MemberAuthor

Static check failures fixed in #21480.

@uranusjr
uranusjr merged commit fded2ca into apache:mainFeb 10, 2022
@uranusjr
uranusjr deleted the decorator-unmap branch February 10, 2022 07:07
ferruzzi pushed a commit to ferruzzi/airflow that referenced this pull request Feb 11, 2022
@jedcunninghamjedcunningham added changelog:skip Changes that should be skipped from the changelog (CI, tests, etc..) area:dynamic-task-mapping AIP-42 labels Feb 28, 2022
@jedcunninghamjedcunningham added this to the Airflow 2.3.0 milestone Apr 26, 2022
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:CLIarea:dynamic-task-mappingAIP-42area:Schedulerincluding HA (high availability) schedulerarea:serializationchangelog:skipChanges that should be skipped from the changelog (CI, tests, etc..)full tests neededWe need to run full set of tests for this PR to mergeprovider:cncf-kubernetesKubernetes (k8s) provider related issues

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants

@uranusjr@ashb@jedcunningham
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Strip utm_, fbclid, gclid, etc. from all links on page\n(function() {\n var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content',\n 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid',\n 'ref', 'ref_src', 'source', 'medium', 'campaign'];\n \n function cleanUrl(url) {\n try {\n var u = new URL(url, window.location.origin);\n var changed = false;\n trackingParams.forEach(function(p) {\n if (u.searchParams.has(p)) {\n u.searchParams.delete(p);\n changed = true;\n }\n });\n return changed ? u.toString() : url;\n } catch (e) {\n return url;\n }\n }\n \n function cleanLinks() {\n document.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n \n cleanLinks();\n \n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1) {\n if (node.tagName === 'A') cleanLinks();\n node.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Remove Tracking Parameters from Links"); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + '
Skip to content

Rewrite decorated task mapping - #21328

Merged
uranusjr merged 6 commits into
apache:mainfrom
astronomer:decorator-unmap
Feb 10, 2022
Merged

Rewrite decorated task mapping#21328
uranusjr merged 6 commits into
apache:mainfrom
astronomer:decorator-unmap

Conversation

@uranusjr

@uranusjruranusjr commented Feb 4, 2022

Copy link
Copy Markdown
Member

Need to get #21210 merged first and rebase this. Done

This rewrites how map() is implemented on _TaskDecorator. The previous implementation incorrectly assumed mapping a task decorator works like mapping a traditional operator, but in fact they have entirely different semantics.

When mapping a traditional operator, FooOperator.map(my_arg=[1, 2, 3]), the argument on FooOperator is mapped. But when mapping a task decorator, my_task.map(my_val=[1, 2, 3]), it’s the argument on the function wrapped by the task object being mapped. Therefore, mapped (and also partial-ed) values for my_val should go into the DecoratedOperator’s op_kwargs (and op_args) arguments instead.

An end-to-end test is provided to validate the simplest case of this. Also fixed a test where the end-to-end test becomes broken after running a DAG (because I need to run a second DAG!)

@boring-cyborgboring-cyborgBot added area:CLI provider:cncf-kubernetes Kubernetes (k8s) provider related issues area:Scheduler including HA (high availability) scheduler area:serialization labels Feb 4, 2022
@uranusjr
uranusjr marked this pull request as ready for review February 4, 2022 14:32
@uranusjr
uranusjrforce-pushed the decorator-unmap branch 6 times, most recently from b52938c to f963a48CompareFebruary 6, 2022 12:11
@ashb

ashb commented Feb 7, 2022

Copy link
Copy Markdown
Member

The serialization needs some work for this approach. Here is an incomplete test:

deftest_mapped_decorator_serde():
fromairflow.models.xcom_argimportXComArgfromairflow.decoratorsimporttaskwithDAG("test-dag", start_date=datetime(2020, 1, 1)) asdag:
task1=BaseOperator(task_id="op1")
xcomarg=XComArg(task1, "test_key")
@task(retry_delay=30)defx(arg1, arg2):
...
real_op=x.partial(arg1=1).map(arg2=xcomarg).operatorserialized=SerializedBaseOperator._serialize(real_op)
assertserialized== {
'_is_dummy': False,
'_is_mapped': True,
'_task_module': 'airflow.decorators.python',
'_task_type': '_PythonDecoratedOperator',
'downstream_task_ids': [],
'partial_kwargs': {
'multiple_outputs': False,
'op_args': [],
'op_kwargs': {
'arg1': [
1,
2,
{"__type": "dict", "__var": {'a': 'b'}},
],
},
'retry_delay': 30,
},
'mapped_kwargs': {
'op_args': [],
# We don't need the __type/__var here!'op_kwargs': {
'arg2': {'__type': 'xcomref', '__var': {'task_id': 'op1', 'key': 'test_key'}},
}
},
'task_id': 'x',
'template_ext': [],
'template_fields': ['op_args', 'op_kwargs'],
# We don't want to include the python source code in the serialized representation# TODO? Where does `retry_delay` go? We need to separate 
}
op=SerializedBaseOperator.deserialize_operator(serialized)
assertisinstance(op, MappedOperator)
assertop.depsisMappedOperator.DEFAULT_DEPS# TODO: add some more asserts here

The serialization

Trying to call real_op.unmap() also throws a key error:

tests/serialization/test_dag_serialization.py:1670: in test_mapped_decorator_serde
real_op.unmap()
airflow/models/baseoperator.py:1882: in unmap
dag._remove_task(self.task_id)
airflow/models/dag.py:2167: in _remove_task
task = self.task_dict.pop(task_id)
E KeyError: 'x'

@uranusjr

Copy link
Copy Markdown
MemberAuthor

Serialisation format is revised. The KeyError is also fixed by implementing slightly smarter logic to merge op_kwargs from mapped_kwargs and partial_kwargs.

@uranusjr

Copy link
Copy Markdown
MemberAuthor

We should also implement some validation when a DAG is parsed to make sure the user pass reasonable values to partial_kwargs. This will be done in a separate PR.

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.

Suggested change
assertdeserialized.retry_delay==timedelta(seconds=30)

(give or take.)

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

At this point retry_delay is still 30 verbatim; it become a timedelta only after unmapped. I’ll add a test elsewhere for this.

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.

It should be deserialized as a timedelta -- it's fine for this value since the constructor for BaseOperator handles it, but other variables might behave differently

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.

start_date for instance.

@uranusjruranusjrFeb 9, 2022

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

I added a test in tests/decorators/test_python.py for this

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

start_date for instance.

Hm some refactoring would be called for to extract the logic out of BaseOperator for reuse. (Note that this affects non-decorator MappedOperator as well.) I think this should be done in a separate PR.

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.

Yeah agreed. We can fix this up later, this unblocks a lot.

Mapping a traditional operator (where arguments go to the operator)
and a task flow operator (where arguments go to the *function*) have
very different semantics, so we need some special code for them.
Previously we were merging partial_kwargs and mapped_kwargs too naively
and did not correctly handle op_args and op_kwargs; those need special
logic due to the mapping semantics of decorated tasks.
Some attributes are removed from serialization to match the format
of the (unmapped) _PythonDecoratedOperator. Some simplication is
implemented to op_kwargs to save some space.
ashb
ashb approved these changes Feb 9, 2022
@github-actions

Copy link
Copy Markdown
Contributor

The PR most likely needs to run full matrix of tests because it modifies parts of the core of Airflow. However, committers might decide to merge it quickly and take the risk. If they don't merge it quickly - please rebase it to the latest main at your convenience, or amend the last commit of the PR, and push it with --force-with-lease.

@github-actionsgithub-actionsBot added the full tests needed We need to run full set of tests for this PR to merge label Feb 9, 2022
@uranusjr

Copy link
Copy Markdown
MemberAuthor

Static check failures fixed in #21480.

@uranusjr
uranusjr merged commit fded2ca into apache:mainFeb 10, 2022
@uranusjr
uranusjr deleted the decorator-unmap branch February 10, 2022 07:07
ferruzzi pushed a commit to ferruzzi/airflow that referenced this pull request Feb 11, 2022
@jedcunninghamjedcunningham added changelog:skip Changes that should be skipped from the changelog (CI, tests, etc..) area:dynamic-task-mapping AIP-42 labels Feb 28, 2022
@jedcunninghamjedcunningham added this to the Airflow 2.3.0 milestone Apr 26, 2022
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:CLIarea:dynamic-task-mappingAIP-42area:Schedulerincluding HA (high availability) schedulerarea:serializationchangelog:skipChanges that should be skipped from the changelog (CI, tests, etc..)full tests neededWe need to run full set of tests for this PR to mergeprovider:cncf-kubernetesKubernetes (k8s) provider related issues

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants

@uranusjr@ashb@jedcunningham
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Auto-enable theater mode on YouTube\n(function() {\n function tryTheater() {\n var btn = document.querySelector('button[aria-label=\"Theater mode\"], ytd-player #player button[title=\"Theater mode\"]');\n if (btn && !btn.classList.contains('activated')) {\n btn.click();\n }\n }\n \n // Try immediately\n tryTheater();\n \n // Try after navigation (SPA)\n var lastUrl = location.href;\n setInterval(function() {\n if (location.href !== lastUrl) {\n lastUrl = location.href;\n setTimeout(tryTheater, 500);\n }\n }, 1000);\n \n // Also try on player load\n var observer = new MutationObserver(tryTheater);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "YouTube Theater Mode Default"); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Rewrite decorated task mapping - #21328

Merged
uranusjr merged 6 commits into
apache:mainfrom
astronomer:decorator-unmap
Feb 10, 2022
Merged

Rewrite decorated task mapping#21328
uranusjr merged 6 commits into
apache:mainfrom
astronomer:decorator-unmap

Conversation

@uranusjr

@uranusjruranusjr commented Feb 4, 2022

Copy link
Copy Markdown
Member

Need to get #21210 merged first and rebase this. Done

This rewrites how map() is implemented on _TaskDecorator. The previous implementation incorrectly assumed mapping a task decorator works like mapping a traditional operator, but in fact they have entirely different semantics.

When mapping a traditional operator, FooOperator.map(my_arg=[1, 2, 3]), the argument on FooOperator is mapped. But when mapping a task decorator, my_task.map(my_val=[1, 2, 3]), it’s the argument on the function wrapped by the task object being mapped. Therefore, mapped (and also partial-ed) values for my_val should go into the DecoratedOperator’s op_kwargs (and op_args) arguments instead.

An end-to-end test is provided to validate the simplest case of this. Also fixed a test where the end-to-end test becomes broken after running a DAG (because I need to run a second DAG!)

@boring-cyborgboring-cyborgBot added area:CLI provider:cncf-kubernetes Kubernetes (k8s) provider related issues area:Scheduler including HA (high availability) scheduler area:serialization labels Feb 4, 2022
@uranusjr
uranusjr marked this pull request as ready for review February 4, 2022 14:32
@uranusjr
uranusjrforce-pushed the decorator-unmap branch 6 times, most recently from b52938c to f963a48CompareFebruary 6, 2022 12:11
@ashb

ashb commented Feb 7, 2022

Copy link
Copy Markdown
Member

The serialization needs some work for this approach. Here is an incomplete test:

deftest_mapped_decorator_serde():
fromairflow.models.xcom_argimportXComArgfromairflow.decoratorsimporttaskwithDAG("test-dag", start_date=datetime(2020, 1, 1)) asdag:
task1=BaseOperator(task_id="op1")
xcomarg=XComArg(task1, "test_key")
@task(retry_delay=30)defx(arg1, arg2):
...
real_op=x.partial(arg1=1).map(arg2=xcomarg).operatorserialized=SerializedBaseOperator._serialize(real_op)
assertserialized== {
'_is_dummy': False,
'_is_mapped': True,
'_task_module': 'airflow.decorators.python',
'_task_type': '_PythonDecoratedOperator',
'downstream_task_ids': [],
'partial_kwargs': {
'multiple_outputs': False,
'op_args': [],
'op_kwargs': {
'arg1': [
1,
2,
{"__type": "dict", "__var": {'a': 'b'}},
],
},
'retry_delay': 30,
},
'mapped_kwargs': {
'op_args': [],
# We don't need the __type/__var here!'op_kwargs': {
'arg2': {'__type': 'xcomref', '__var': {'task_id': 'op1', 'key': 'test_key'}},
}
},
'task_id': 'x',
'template_ext': [],
'template_fields': ['op_args', 'op_kwargs'],
# We don't want to include the python source code in the serialized representation# TODO? Where does `retry_delay` go? We need to separate 
}
op=SerializedBaseOperator.deserialize_operator(serialized)
assertisinstance(op, MappedOperator)
assertop.depsisMappedOperator.DEFAULT_DEPS# TODO: add some more asserts here

The serialization

Trying to call real_op.unmap() also throws a key error:

tests/serialization/test_dag_serialization.py:1670: in test_mapped_decorator_serde
real_op.unmap()
airflow/models/baseoperator.py:1882: in unmap
dag._remove_task(self.task_id)
airflow/models/dag.py:2167: in _remove_task
task = self.task_dict.pop(task_id)
E KeyError: 'x'

@uranusjr

Copy link
Copy Markdown
MemberAuthor

Serialisation format is revised. The KeyError is also fixed by implementing slightly smarter logic to merge op_kwargs from mapped_kwargs and partial_kwargs.

@uranusjr

Copy link
Copy Markdown
MemberAuthor

We should also implement some validation when a DAG is parsed to make sure the user pass reasonable values to partial_kwargs. This will be done in a separate PR.

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.

Suggested change
assertdeserialized.retry_delay==timedelta(seconds=30)

(give or take.)

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

At this point retry_delay is still 30 verbatim; it become a timedelta only after unmapped. I’ll add a test elsewhere for this.

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.

It should be deserialized as a timedelta -- it's fine for this value since the constructor for BaseOperator handles it, but other variables might behave differently

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.

start_date for instance.

@uranusjruranusjrFeb 9, 2022

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

I added a test in tests/decorators/test_python.py for this

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

start_date for instance.

Hm some refactoring would be called for to extract the logic out of BaseOperator for reuse. (Note that this affects non-decorator MappedOperator as well.) I think this should be done in a separate PR.

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.

Yeah agreed. We can fix this up later, this unblocks a lot.

Mapping a traditional operator (where arguments go to the operator)
and a task flow operator (where arguments go to the *function*) have
very different semantics, so we need some special code for them.
Previously we were merging partial_kwargs and mapped_kwargs too naively
and did not correctly handle op_args and op_kwargs; those need special
logic due to the mapping semantics of decorated tasks.
Some attributes are removed from serialization to match the format
of the (unmapped) _PythonDecoratedOperator. Some simplication is
implemented to op_kwargs to save some space.
ashb
ashb approved these changes Feb 9, 2022
@github-actions

Copy link
Copy Markdown
Contributor

The PR most likely needs to run full matrix of tests because it modifies parts of the core of Airflow. However, committers might decide to merge it quickly and take the risk. If they don't merge it quickly - please rebase it to the latest main at your convenience, or amend the last commit of the PR, and push it with --force-with-lease.

@github-actionsgithub-actionsBot added the full tests needed We need to run full set of tests for this PR to merge label Feb 9, 2022
@uranusjr

Copy link
Copy Markdown
MemberAuthor

Static check failures fixed in #21480.

@uranusjr
uranusjr merged commit fded2ca into apache:mainFeb 10, 2022
@uranusjr
uranusjr deleted the decorator-unmap branch February 10, 2022 07:07
ferruzzi pushed a commit to ferruzzi/airflow that referenced this pull request Feb 11, 2022
@jedcunninghamjedcunningham added changelog:skip Changes that should be skipped from the changelog (CI, tests, etc..) area:dynamic-task-mapping AIP-42 labels Feb 28, 2022
@jedcunninghamjedcunningham added this to the Airflow 2.3.0 milestone Apr 26, 2022
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:CLIarea:dynamic-task-mappingAIP-42area:Schedulerincluding HA (high availability) schedulerarea:serializationchangelog:skipChanges that should be skipped from the changelog (CI, tests, etc..)full tests neededWe need to run full set of tests for this PR to mergeprovider:cncf-kubernetesKubernetes (k8s) provider related issues

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants

@uranusjr@ashb@jedcunningham
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Remove or un-stick sticky/fixed headers that block content\n(function() {\n function unstick() {\n document.querySelectorAll('header, nav, [role=\"banner\"], .header, .navbar, .sticky, .fixed-top, [style*=\"position: fixed\"], [style*=\"position:sticky\"]').forEach(function(el) {\n if (el.style.position === 'fixed' || el.style.position === 'sticky' || \n getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') {\n el.style.position = 'static';\n el.style.top = 'auto';\n el.style.zIndex = 'auto';\n }\n });\n }\n \n unstick();\n \n var observer = new MutationObserver(unstick);\n observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] });\n})();", "Kill Sticky Headers"); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Rewrite decorated task mapping - #21328

Merged
uranusjr merged 6 commits into
apache:mainfrom
astronomer:decorator-unmap
Feb 10, 2022
Merged

Rewrite decorated task mapping#21328
uranusjr merged 6 commits into
apache:mainfrom
astronomer:decorator-unmap

Conversation

@uranusjr

@uranusjruranusjr commented Feb 4, 2022

Copy link
Copy Markdown
Member

Need to get #21210 merged first and rebase this. Done

This rewrites how map() is implemented on _TaskDecorator. The previous implementation incorrectly assumed mapping a task decorator works like mapping a traditional operator, but in fact they have entirely different semantics.

When mapping a traditional operator, FooOperator.map(my_arg=[1, 2, 3]), the argument on FooOperator is mapped. But when mapping a task decorator, my_task.map(my_val=[1, 2, 3]), it’s the argument on the function wrapped by the task object being mapped. Therefore, mapped (and also partial-ed) values for my_val should go into the DecoratedOperator’s op_kwargs (and op_args) arguments instead.

An end-to-end test is provided to validate the simplest case of this. Also fixed a test where the end-to-end test becomes broken after running a DAG (because I need to run a second DAG!)

@boring-cyborgboring-cyborgBot added area:CLI provider:cncf-kubernetes Kubernetes (k8s) provider related issues area:Scheduler including HA (high availability) scheduler area:serialization labels Feb 4, 2022
@uranusjr
uranusjr marked this pull request as ready for review February 4, 2022 14:32
@uranusjr
uranusjrforce-pushed the decorator-unmap branch 6 times, most recently from b52938c to f963a48CompareFebruary 6, 2022 12:11
@ashb

ashb commented Feb 7, 2022

Copy link
Copy Markdown
Member

The serialization needs some work for this approach. Here is an incomplete test:

deftest_mapped_decorator_serde():
fromairflow.models.xcom_argimportXComArgfromairflow.decoratorsimporttaskwithDAG("test-dag", start_date=datetime(2020, 1, 1)) asdag:
task1=BaseOperator(task_id="op1")
xcomarg=XComArg(task1, "test_key")
@task(retry_delay=30)defx(arg1, arg2):
...
real_op=x.partial(arg1=1).map(arg2=xcomarg).operatorserialized=SerializedBaseOperator._serialize(real_op)
assertserialized== {
'_is_dummy': False,
'_is_mapped': True,
'_task_module': 'airflow.decorators.python',
'_task_type': '_PythonDecoratedOperator',
'downstream_task_ids': [],
'partial_kwargs': {
'multiple_outputs': False,
'op_args': [],
'op_kwargs': {
'arg1': [
1,
2,
{"__type": "dict", "__var": {'a': 'b'}},
],
},
'retry_delay': 30,
},
'mapped_kwargs': {
'op_args': [],
# We don't need the __type/__var here!'op_kwargs': {
'arg2': {'__type': 'xcomref', '__var': {'task_id': 'op1', 'key': 'test_key'}},
}
},
'task_id': 'x',
'template_ext': [],
'template_fields': ['op_args', 'op_kwargs'],
# We don't want to include the python source code in the serialized representation# TODO? Where does `retry_delay` go? We need to separate 
}
op=SerializedBaseOperator.deserialize_operator(serialized)
assertisinstance(op, MappedOperator)
assertop.depsisMappedOperator.DEFAULT_DEPS# TODO: add some more asserts here

The serialization

Trying to call real_op.unmap() also throws a key error:

tests/serialization/test_dag_serialization.py:1670: in test_mapped_decorator_serde
real_op.unmap()
airflow/models/baseoperator.py:1882: in unmap
dag._remove_task(self.task_id)
airflow/models/dag.py:2167: in _remove_task
task = self.task_dict.pop(task_id)
E KeyError: 'x'

@uranusjr

Copy link
Copy Markdown
MemberAuthor

Serialisation format is revised. The KeyError is also fixed by implementing slightly smarter logic to merge op_kwargs from mapped_kwargs and partial_kwargs.

@uranusjr

Copy link
Copy Markdown
MemberAuthor

We should also implement some validation when a DAG is parsed to make sure the user pass reasonable values to partial_kwargs. This will be done in a separate PR.

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.

Suggested change
assertdeserialized.retry_delay==timedelta(seconds=30)

(give or take.)

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

At this point retry_delay is still 30 verbatim; it become a timedelta only after unmapped. I’ll add a test elsewhere for this.

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.

It should be deserialized as a timedelta -- it's fine for this value since the constructor for BaseOperator handles it, but other variables might behave differently

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.

start_date for instance.

@uranusjruranusjrFeb 9, 2022

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

I added a test in tests/decorators/test_python.py for this

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

start_date for instance.

Hm some refactoring would be called for to extract the logic out of BaseOperator for reuse. (Note that this affects non-decorator MappedOperator as well.) I think this should be done in a separate PR.

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.

Yeah agreed. We can fix this up later, this unblocks a lot.

Mapping a traditional operator (where arguments go to the operator)
and a task flow operator (where arguments go to the *function*) have
very different semantics, so we need some special code for them.
Previously we were merging partial_kwargs and mapped_kwargs too naively
and did not correctly handle op_args and op_kwargs; those need special
logic due to the mapping semantics of decorated tasks.
Some attributes are removed from serialization to match the format
of the (unmapped) _PythonDecoratedOperator. Some simplication is
implemented to op_kwargs to save some space.
ashb
ashb approved these changes Feb 9, 2022
@github-actions

Copy link
Copy Markdown
Contributor

The PR most likely needs to run full matrix of tests because it modifies parts of the core of Airflow. However, committers might decide to merge it quickly and take the risk. If they don't merge it quickly - please rebase it to the latest main at your convenience, or amend the last commit of the PR, and push it with --force-with-lease.

@github-actionsgithub-actionsBot added the full tests needed We need to run full set of tests for this PR to merge label Feb 9, 2022
@uranusjr

Copy link
Copy Markdown
MemberAuthor

Static check failures fixed in #21480.

@uranusjr
uranusjr merged commit fded2ca into apache:mainFeb 10, 2022
@uranusjr
uranusjr deleted the decorator-unmap branch February 10, 2022 07:07
ferruzzi pushed a commit to ferruzzi/airflow that referenced this pull request Feb 11, 2022
@jedcunninghamjedcunningham added changelog:skip Changes that should be skipped from the changelog (CI, tests, etc..) area:dynamic-task-mapping AIP-42 labels Feb 28, 2022
@jedcunninghamjedcunningham added this to the Airflow 2.3.0 milestone Apr 26, 2022
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:CLIarea:dynamic-task-mappingAIP-42area:Schedulerincluding HA (high availability) schedulerarea:serializationchangelog:skipChanges that should be skipped from the changelog (CI, tests, etc..)full tests neededWe need to run full set of tests for this PR to mergeprovider:cncf-kubernetesKubernetes (k8s) provider related issues

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants

@uranusjr@ashb@jedcunningham
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Universal Dark Mode - works on any site\n(function() {\n var enabled = true;\n \n function applyDarkMode() {\n if (!enabled) return;\n \n // Create style element if it doesn't exist\n var style = document.getElementById('universal-dark-mode-style');\n if (!style) {\n style = document.createElement('style');\n style.id = 'universal-dark-mode-style';\n document.head.appendChild(style);\n }\n \n // Dark mode CSS - inverts colors but preserves images/video\n style.textContent = '\n /* Invert everything except media */\n html {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #1a1a2e !important;\n }\n \n /* Restore images, videos, iframes, canvas */\n img, video, iframe, canvas, svg, picture, [style*=\"background-image\"] {\n filter: invert(1) hue-rotate(180deg) !important;\n }\n \n /* Preserve specific elements that should not be inverted */\n .no-dark-mode, .no-dark-mode *,\n [data-theme=\"light\"], [data-theme=\"light\"],\n .ace_editor, .ace_editor *,\n .CodeMirror, .CodeMirror *,\n .monaco-editor, .monaco-editor *,\n .markdown-body pre, .markdown-body pre *,\n .highlight, .highlight *,\n pre code, pre code * {\n filter: none !important;\n }\n \n /* Fix common UI elements */\n .modal, .popup, .dropdown-menu, .tooltip, .popover {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #2d2d44 !important;\n border-color: #444 !important;\n }\n \n /* Scrollbars */\n ::-webkit-scrollbar { background: #1a1a2e !important; }\n ::-webkit-scrollbar-thumb { background: #444 !important; }\n ::-webkit-scrollbar-thumb:hover { background: #555 !important; }\n \n /* Selection */\n ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ';\n }\n \n function removeDarkMode() {\n var style = document.getElementById('universal-dark-mode-style');\n if (style) style.remove();\n }\n \n // Toggle with Alt+Shift+D\n document.addEventListener('keydown', function(e) {\n if (e.altKey && e.shiftKey && e.key === 'D') {\n e.preventDefault();\n enabled = !enabled;\n if (enabled) {\n applyDarkMode();\n console.log('[Universal Dark Mode] Enabled');\n } else {\n removeDarkMode();\n console.log('[Universal Dark Mode] Disabled');\n }\n }\n });\n \n // Apply on load\n applyDarkMode();\n \n // Re-apply on dynamic content\n var observer = new MutationObserver(function(mutations) {\n if (enabled && !document.getElementById('universal-dark-mode-style')) {\n applyDarkMode();\n }\n });\n observer.observe(document.head, { childList: true });\n \n console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle');\n})();", "Universal Dark Mode"); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })();
Skip to content

Rewrite decorated task mapping - #21328

Merged
uranusjr merged 6 commits into
apache:mainfrom
astronomer:decorator-unmap
Feb 10, 2022
Merged

Rewrite decorated task mapping#21328
uranusjr merged 6 commits into
apache:mainfrom
astronomer:decorator-unmap

Conversation

@uranusjr

@uranusjruranusjr commented Feb 4, 2022

Copy link
Copy Markdown
Member

Need to get #21210 merged first and rebase this. Done

This rewrites how map() is implemented on _TaskDecorator. The previous implementation incorrectly assumed mapping a task decorator works like mapping a traditional operator, but in fact they have entirely different semantics.

When mapping a traditional operator, FooOperator.map(my_arg=[1, 2, 3]), the argument on FooOperator is mapped. But when mapping a task decorator, my_task.map(my_val=[1, 2, 3]), it’s the argument on the function wrapped by the task object being mapped. Therefore, mapped (and also partial-ed) values for my_val should go into the DecoratedOperator’s op_kwargs (and op_args) arguments instead.

An end-to-end test is provided to validate the simplest case of this. Also fixed a test where the end-to-end test becomes broken after running a DAG (because I need to run a second DAG!)

@boring-cyborgboring-cyborgBot added area:CLI provider:cncf-kubernetes Kubernetes (k8s) provider related issues area:Scheduler including HA (high availability) scheduler area:serialization labels Feb 4, 2022
@uranusjr
uranusjr marked this pull request as ready for review February 4, 2022 14:32
@uranusjr
uranusjrforce-pushed the decorator-unmap branch 6 times, most recently from b52938c to f963a48CompareFebruary 6, 2022 12:11
@ashb

ashb commented Feb 7, 2022

Copy link
Copy Markdown
Member

The serialization needs some work for this approach. Here is an incomplete test:

deftest_mapped_decorator_serde():
fromairflow.models.xcom_argimportXComArgfromairflow.decoratorsimporttaskwithDAG("test-dag", start_date=datetime(2020, 1, 1)) asdag:
task1=BaseOperator(task_id="op1")
xcomarg=XComArg(task1, "test_key")
@task(retry_delay=30)defx(arg1, arg2):
...
real_op=x.partial(arg1=1).map(arg2=xcomarg).operatorserialized=SerializedBaseOperator._serialize(real_op)
assertserialized== {
'_is_dummy': False,
'_is_mapped': True,
'_task_module': 'airflow.decorators.python',
'_task_type': '_PythonDecoratedOperator',
'downstream_task_ids': [],
'partial_kwargs': {
'multiple_outputs': False,
'op_args': [],
'op_kwargs': {
'arg1': [
1,
2,
{"__type": "dict", "__var": {'a': 'b'}},
],
},
'retry_delay': 30,
},
'mapped_kwargs': {
'op_args': [],
# We don't need the __type/__var here!'op_kwargs': {
'arg2': {'__type': 'xcomref', '__var': {'task_id': 'op1', 'key': 'test_key'}},
}
},
'task_id': 'x',
'template_ext': [],
'template_fields': ['op_args', 'op_kwargs'],
# We don't want to include the python source code in the serialized representation# TODO? Where does `retry_delay` go? We need to separate 
}
op=SerializedBaseOperator.deserialize_operator(serialized)
assertisinstance(op, MappedOperator)
assertop.depsisMappedOperator.DEFAULT_DEPS# TODO: add some more asserts here

The serialization

Trying to call real_op.unmap() also throws a key error:

tests/serialization/test_dag_serialization.py:1670: in test_mapped_decorator_serde
real_op.unmap()
airflow/models/baseoperator.py:1882: in unmap
dag._remove_task(self.task_id)
airflow/models/dag.py:2167: in _remove_task
task = self.task_dict.pop(task_id)
E KeyError: 'x'

@uranusjr

Copy link
Copy Markdown
MemberAuthor

Serialisation format is revised. The KeyError is also fixed by implementing slightly smarter logic to merge op_kwargs from mapped_kwargs and partial_kwargs.

@uranusjr

Copy link
Copy Markdown
MemberAuthor

We should also implement some validation when a DAG is parsed to make sure the user pass reasonable values to partial_kwargs. This will be done in a separate PR.

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.

Suggested change
assertdeserialized.retry_delay==timedelta(seconds=30)

(give or take.)

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

At this point retry_delay is still 30 verbatim; it become a timedelta only after unmapped. I’ll add a test elsewhere for this.

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.

It should be deserialized as a timedelta -- it's fine for this value since the constructor for BaseOperator handles it, but other variables might behave differently

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.

start_date for instance.

@uranusjruranusjrFeb 9, 2022

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

I added a test in tests/decorators/test_python.py for this

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

start_date for instance.

Hm some refactoring would be called for to extract the logic out of BaseOperator for reuse. (Note that this affects non-decorator MappedOperator as well.) I think this should be done in a separate PR.

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.

Yeah agreed. We can fix this up later, this unblocks a lot.

Mapping a traditional operator (where arguments go to the operator)
and a task flow operator (where arguments go to the *function*) have
very different semantics, so we need some special code for them.
Previously we were merging partial_kwargs and mapped_kwargs too naively
and did not correctly handle op_args and op_kwargs; those need special
logic due to the mapping semantics of decorated tasks.
Some attributes are removed from serialization to match the format
of the (unmapped) _PythonDecoratedOperator. Some simplication is
implemented to op_kwargs to save some space.
ashb
ashb approved these changes Feb 9, 2022
@github-actions

Copy link
Copy Markdown
Contributor

The PR most likely needs to run full matrix of tests because it modifies parts of the core of Airflow. However, committers might decide to merge it quickly and take the risk. If they don't merge it quickly - please rebase it to the latest main at your convenience, or amend the last commit of the PR, and push it with --force-with-lease.

@github-actionsgithub-actionsBot added the full tests needed We need to run full set of tests for this PR to merge label Feb 9, 2022
@uranusjr

Copy link
Copy Markdown
MemberAuthor

Static check failures fixed in #21480.

@uranusjr
uranusjr merged commit fded2ca into apache:mainFeb 10, 2022
@uranusjr
uranusjr deleted the decorator-unmap branch February 10, 2022 07:07
ferruzzi pushed a commit to ferruzzi/airflow that referenced this pull request Feb 11, 2022
@jedcunninghamjedcunningham added changelog:skip Changes that should be skipped from the changelog (CI, tests, etc..) area:dynamic-task-mapping AIP-42 labels Feb 28, 2022
@jedcunninghamjedcunningham added this to the Airflow 2.3.0 milestone Apr 26, 2022
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:CLIarea:dynamic-task-mappingAIP-42area:Schedulerincluding HA (high availability) schedulerarea:serializationchangelog:skipChanges that should be skipped from the changelog (CI, tests, etc..)full tests neededWe need to run full set of tests for this PR to mergeprovider:cncf-kubernetesKubernetes (k8s) provider related issues

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants

@uranusjr@ashb@jedcunningham