AIP-44 Support TaskInstance serialization/deserialization. - #29355

Closed
mhenc wants to merge 7 commits into
apache:mainfrom
mhenc:ti_serialize
Closed

AIP-44 Support TaskInstance serialization/deserialization.#29355
mhenc wants to merge 7 commits into
apache:mainfrom
mhenc:ti_serialize

Conversation

@mhenc

@mhencmhenc commented Feb 3, 2023

Copy link
Copy Markdown
Contributor

In AIP-44 we need to send the TaskInstance object to worker for execution. Currently this object type doesn't support serialization(only SimpleTaskInstance does).

This change is mainly about new way to construct TaskInstance from_dict - which requires changing the constructor into from_task and refactoring all usages.

Note that deserialized TaskInstance will not have ORM state (so will not be able to refresh state automatically from DB etc), but this should be fine as it will be used on client side of Internal API.

related: #29320

@boring-cyborgboring-cyborgBot added area:API Airflow's REST/HTTP API area:CLI area:core-operators provider:cncf-kubernetes Kubernetes (k8s) provider related issues area:Scheduler including HA (high availability) scheduler area:serialization area:webserver Webserver related Issues labels Feb 3, 2023
@mhenc
mhenc marked this pull request as ready for review February 3, 2023 14:16
@mhenc

mhenc commented Feb 3, 2023

Copy link
Copy Markdown
ContributorAuthor

cc: @potiuk@vincbeck

@mhencmhenc changed the title AIP-44 Support TaskInstance serialization.AIP-44 Support TaskInstance serialization/deserialization.Feb 3, 2023

@vincbeckvincbeck left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Nice!

Comment threadairflow/models/taskinstance.py Outdated

@kosteevkosteev left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

In order to achieve the goal of constructing TaskInstance objects from "task" and "dict"
could we have constructor of TaskInstance extended to support new parameter, e.g. task_dict like this:

def __init__(
self,
task: Operator,
execution_date: datetime | None = None,
run_id: str | None = None,
state: str | None = None,
map_index: int = -1,
task_dict: dict = None,
):
if task_dict is not None:
# init from task_dict
return
# old implementation

and then TaskInstance(task=t), TaskInstance(task_dict={...})

What do you think about this approach?

Comment threadairflow/models/taskinstance.py Outdated
@ashb

ashb commented Feb 8, 2023

Copy link
Copy Markdown
Member

@bolkedebruin could you take a look? I know you've recently overhauled all of the serialisation code

1 similar comment
@ashb

ashb commented Feb 8, 2023

Copy link
Copy Markdown
Member

@bolkedebruin could you take a look? I know you've recently overhauled all of the serialisation code

@bolkedebruinbolkedebruin left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Overall, great that this is happening! I would like to bring it more inline with the new serialization/deserialization code as I have been working on making it into 'one' (merging XCOM serialization and DAG serialization).

To seduce you: the new serialization code is an order of a magnitude faster than the old one (improvements over 50%) and can deal with versioning, while the old one cannot. The migration isn't fully completed yet (I am working on DAGs, but for obvious reasons that will take awhile), however that should not be in the way of making use of the new code (serde.py).

For your code this means moving away from to_dict/from_dict to the API of serialize() / deserialize(value, version) and setting __version__: ClassVar[int] = X in the TaskInstance class. This is pretty close to what you have now. serialization/deserialization should be called from serde and not directly. Alternatively you can add a serializer/deserializer to airflow.serialization.serializers.

I can help if needed.

Comment threadairflow/models/taskinstance.py Outdated
Comment threadairflow/models/taskinstance.py Outdated

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Here you are extending legacy code, i suggest using the more generic serialization code from serde

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

We rely on serialized_objects in InternalAPI in general
https://github.com/apache/airflow/blob/main/airflow/api_internal/internal_api_call.py#L107
Do you think we could switch to serde now? Is it compatible with what serialize_objects offer?

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Given that you are going to rely on a lot of serialization and the new serializer/deserializer is significantly faster and more future proof (versioning) I think everything that does currently not have a schema (that's everything except a DAG*) should switch.

I am willing to help out to ease the migration if required.

  • I am working on DAG serialization/deserialization but untangling how it is done now and to improve the structure is taking time especially with all the edge cases.

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

Great. I think we will migrate to the new serializer/deserializer in this case, but probably outside of this PR.
If you believe it would better to migrate first, then I can revert this change and get back to it when using new way.
WDYT?

@bolkedebruinbolkedebruinFeb 9, 2023

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Am I incorrect in thinking that migration is a 2 line change in internal_api_call and you are not relying on any of the other serialized_objects (basically DAG)? If so then I would say do not add technical dept and migrate now. This allows us to call serialized_objects as stale and soon to be deprecated.

Otherwise, keep it and and add it to the todo of AIP-44?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

yes, we need to make change there and server-side:
https://github.com/apache/airflow/blob/main/airflow/api_internal/endpoints/rpc_api_endpoint.py#L76

There are more methods with internal_api_call decorator. I did a quick check and I see that (beside primitives) we already need serialization to Dag,DagRun, BaseXCom, CallbackRequest (and probably more soon)

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.

Question. (and really sorry I have not looked at it earlier) and also follow up on #29513 as well.

Are we REALLY sure we need to serialize the whole TaskInstance (and other) objects to be passed via internal_api_call ?

For me this is an indication that we either have too narrow of a scope for an @internal_api call (generally speaking the whole internal_api_call should span the whole DB transaction. And since we are trying to pass an ORM object (TaskInstance, DAGRun etc.) it means that that object must have been retrieved before within a transaction. So it means that our internal_api_call should wrap the retrieval as well. Which might simply mean that we need to do some refactoring and add extract new methods (and then decorate them).

That's of course a general statement and approach and there might be cases that this require a bit deeper refactoring.

Which methods are affected @mhenc (besides the #29513 one) ? maybe we can look toghether and figure out approach for all of them ?

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.

For example in #29513 (comment) I proposed the solution that would avoid TaskInstance serialization altogether. I am reasonable convinced, that similar approach can be done for all ORM objects of ours and that we do not need to serialize any of them (in which case the whole PR might not be needed).

@mhencmhencFeb 21, 2023

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

But this works only for client -> Internal API side.

What about Internal API -> client, e.g. worker. We need to have TaskInstance object in worker to run the task, e.g.
https://github.com/apache/airflow/blob/main/airflow/cli/commands/task_command.py#L187

unless of course we are able to refactor it completely

@potiukpotiukFeb 21, 2023

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.

That's precisely what I am refactoring now (with my POC)

@bolkedebruin

bolkedebruin commented Feb 8, 2023

Copy link
Copy Markdown
Contributor

In order to achieve the goal of constructing TaskInstance objects from "task" and "dict" could we have constructor of TaskInstance extended to support new parameter, e.g. task_dict like this:

def init(
self,
task: Operator,
execution_date: datetime | None = None,
run_id: str | None = None,
state: str | None = None,
map_index: int = -1,
task_dict: dict = None,
):
if task_dict is not None:
# init from task_dict
return

# old implementation

and then TaskInstance(task=t), TaskInstance(task_dict={...})

What do you think about this approach?

I don't think that makes sense, a staticmethod or classmethod deserialize should cover this and used this way at other locations.

It could be as simple as this:

from airflow.serialization.serde import deserialize
ti = deserialize(serialized_ti)

@bolkedebruinbolkedebruin left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Added some additional comments

Comment threadairflow/models/taskinstance.py Outdated
Comment threadairflow/models/taskinstance.py Outdated
Comment threadairflow/models/taskinstance.py Outdated
@mhenc
mhenc requested review from bolkedebruin and kosteev and removed request for kosteevFebruary 9, 2023 13:38
@mhenc
mhenc requested a review from uranusjrFebruary 9, 2023 13:39

@bolkedebruinbolkedebruin left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Almost there :-)

Comment threadairflow/models/taskinstance.py Outdated
Comment threadairflow/models/taskinstance.py Outdated

@bolkedebruinbolkedebruinFeb 9, 2023

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Am I incorrect in thinking that migration is a 2 line change in internal_api_call and you are not relying on any of the other serialized_objects (basically DAG)? If so then I would say do not add technical dept and migrate now. This allows us to call serialized_objects as stale and soon to be deprecated.

Otherwise, keep it and and add it to the todo of AIP-44?

@mhenc

Copy link
Copy Markdown
ContributorAuthor

Closing in favor of #30282

@mhencmhenc closed this Mar 24, 2023
@mhenc
mhenc deleted the ti_serialize branch August 30, 2023 18:48
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:APIAirflow's REST/HTTP APIarea:CLIarea:Schedulerincluding HA (high availability) schedulerarea:serializationarea:webserverWebserver related Issuesprovider:cncf-kubernetesKubernetes (k8s) provider related issues

Projects

None yet

Development

Successfully merging this pull request may close these issues.

7 participants

@mhenc@ashb@bolkedebruin@potiuk@uranusjr@kosteev@vincbeck
, '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

AIP-44 Support TaskInstance serialization/deserialization. - #29355

Closed
mhenc wants to merge 7 commits into
apache:mainfrom
mhenc:ti_serialize
Closed

AIP-44 Support TaskInstance serialization/deserialization.#29355
mhenc wants to merge 7 commits into
apache:mainfrom
mhenc:ti_serialize

Conversation

@mhenc

@mhencmhenc commented Feb 3, 2023

Copy link
Copy Markdown
Contributor

In AIP-44 we need to send the TaskInstance object to worker for execution. Currently this object type doesn't support serialization(only SimpleTaskInstance does).

This change is mainly about new way to construct TaskInstance from_dict - which requires changing the constructor into from_task and refactoring all usages.

Note that deserialized TaskInstance will not have ORM state (so will not be able to refresh state automatically from DB etc), but this should be fine as it will be used on client side of Internal API.

related: #29320

@boring-cyborgboring-cyborgBot added area:API Airflow's REST/HTTP API area:CLI area:core-operators provider:cncf-kubernetes Kubernetes (k8s) provider related issues area:Scheduler including HA (high availability) scheduler area:serialization area:webserver Webserver related Issues labels Feb 3, 2023
@mhenc
mhenc marked this pull request as ready for review February 3, 2023 14:16
@mhenc

mhenc commented Feb 3, 2023

Copy link
Copy Markdown
ContributorAuthor

cc: @potiuk@vincbeck

@mhencmhenc changed the title AIP-44 Support TaskInstance serialization.AIP-44 Support TaskInstance serialization/deserialization.Feb 3, 2023

@vincbeckvincbeck left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Nice!

Comment threadairflow/models/taskinstance.py Outdated

@kosteevkosteev left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

In order to achieve the goal of constructing TaskInstance objects from "task" and "dict"
could we have constructor of TaskInstance extended to support new parameter, e.g. task_dict like this:

def __init__(
self,
task: Operator,
execution_date: datetime | None = None,
run_id: str | None = None,
state: str | None = None,
map_index: int = -1,
task_dict: dict = None,
):
if task_dict is not None:
# init from task_dict
return
# old implementation

and then TaskInstance(task=t), TaskInstance(task_dict={...})

What do you think about this approach?

Comment threadairflow/models/taskinstance.py Outdated
@ashb

ashb commented Feb 8, 2023

Copy link
Copy Markdown
Member

@bolkedebruin could you take a look? I know you've recently overhauled all of the serialisation code

1 similar comment
@ashb

ashb commented Feb 8, 2023

Copy link
Copy Markdown
Member

@bolkedebruin could you take a look? I know you've recently overhauled all of the serialisation code

@bolkedebruinbolkedebruin left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Overall, great that this is happening! I would like to bring it more inline with the new serialization/deserialization code as I have been working on making it into 'one' (merging XCOM serialization and DAG serialization).

To seduce you: the new serialization code is an order of a magnitude faster than the old one (improvements over 50%) and can deal with versioning, while the old one cannot. The migration isn't fully completed yet (I am working on DAGs, but for obvious reasons that will take awhile), however that should not be in the way of making use of the new code (serde.py).

For your code this means moving away from to_dict/from_dict to the API of serialize() / deserialize(value, version) and setting __version__: ClassVar[int] = X in the TaskInstance class. This is pretty close to what you have now. serialization/deserialization should be called from serde and not directly. Alternatively you can add a serializer/deserializer to airflow.serialization.serializers.

I can help if needed.

Comment threadairflow/models/taskinstance.py Outdated
Comment threadairflow/models/taskinstance.py Outdated

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Here you are extending legacy code, i suggest using the more generic serialization code from serde

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

We rely on serialized_objects in InternalAPI in general
https://github.com/apache/airflow/blob/main/airflow/api_internal/internal_api_call.py#L107
Do you think we could switch to serde now? Is it compatible with what serialize_objects offer?

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Given that you are going to rely on a lot of serialization and the new serializer/deserializer is significantly faster and more future proof (versioning) I think everything that does currently not have a schema (that's everything except a DAG*) should switch.

I am willing to help out to ease the migration if required.

  • I am working on DAG serialization/deserialization but untangling how it is done now and to improve the structure is taking time especially with all the edge cases.

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

Great. I think we will migrate to the new serializer/deserializer in this case, but probably outside of this PR.
If you believe it would better to migrate first, then I can revert this change and get back to it when using new way.
WDYT?

@bolkedebruinbolkedebruinFeb 9, 2023

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Am I incorrect in thinking that migration is a 2 line change in internal_api_call and you are not relying on any of the other serialized_objects (basically DAG)? If so then I would say do not add technical dept and migrate now. This allows us to call serialized_objects as stale and soon to be deprecated.

Otherwise, keep it and and add it to the todo of AIP-44?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

yes, we need to make change there and server-side:
https://github.com/apache/airflow/blob/main/airflow/api_internal/endpoints/rpc_api_endpoint.py#L76

There are more methods with internal_api_call decorator. I did a quick check and I see that (beside primitives) we already need serialization to Dag,DagRun, BaseXCom, CallbackRequest (and probably more soon)

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.

Question. (and really sorry I have not looked at it earlier) and also follow up on #29513 as well.

Are we REALLY sure we need to serialize the whole TaskInstance (and other) objects to be passed via internal_api_call ?

For me this is an indication that we either have too narrow of a scope for an @internal_api call (generally speaking the whole internal_api_call should span the whole DB transaction. And since we are trying to pass an ORM object (TaskInstance, DAGRun etc.) it means that that object must have been retrieved before within a transaction. So it means that our internal_api_call should wrap the retrieval as well. Which might simply mean that we need to do some refactoring and add extract new methods (and then decorate them).

That's of course a general statement and approach and there might be cases that this require a bit deeper refactoring.

Which methods are affected @mhenc (besides the #29513 one) ? maybe we can look toghether and figure out approach for all of them ?

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.

For example in #29513 (comment) I proposed the solution that would avoid TaskInstance serialization altogether. I am reasonable convinced, that similar approach can be done for all ORM objects of ours and that we do not need to serialize any of them (in which case the whole PR might not be needed).

@mhencmhencFeb 21, 2023

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

But this works only for client -> Internal API side.

What about Internal API -> client, e.g. worker. We need to have TaskInstance object in worker to run the task, e.g.
https://github.com/apache/airflow/blob/main/airflow/cli/commands/task_command.py#L187

unless of course we are able to refactor it completely

@potiukpotiukFeb 21, 2023

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.

That's precisely what I am refactoring now (with my POC)

@bolkedebruin

bolkedebruin commented Feb 8, 2023

Copy link
Copy Markdown
Contributor

In order to achieve the goal of constructing TaskInstance objects from "task" and "dict" could we have constructor of TaskInstance extended to support new parameter, e.g. task_dict like this:

def init(
self,
task: Operator,
execution_date: datetime | None = None,
run_id: str | None = None,
state: str | None = None,
map_index: int = -1,
task_dict: dict = None,
):
if task_dict is not None:
# init from task_dict
return

# old implementation

and then TaskInstance(task=t), TaskInstance(task_dict={...})

What do you think about this approach?

I don't think that makes sense, a staticmethod or classmethod deserialize should cover this and used this way at other locations.

It could be as simple as this:

from airflow.serialization.serde import deserialize
ti = deserialize(serialized_ti)

@bolkedebruinbolkedebruin left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Added some additional comments

Comment threadairflow/models/taskinstance.py Outdated
Comment threadairflow/models/taskinstance.py Outdated
Comment threadairflow/models/taskinstance.py Outdated
@mhenc
mhenc requested review from bolkedebruin and kosteev and removed request for kosteevFebruary 9, 2023 13:38
@mhenc
mhenc requested a review from uranusjrFebruary 9, 2023 13:39

@bolkedebruinbolkedebruin left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Almost there :-)

Comment threadairflow/models/taskinstance.py Outdated
Comment threadairflow/models/taskinstance.py Outdated

@bolkedebruinbolkedebruinFeb 9, 2023

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Am I incorrect in thinking that migration is a 2 line change in internal_api_call and you are not relying on any of the other serialized_objects (basically DAG)? If so then I would say do not add technical dept and migrate now. This allows us to call serialized_objects as stale and soon to be deprecated.

Otherwise, keep it and and add it to the todo of AIP-44?

@mhenc

Copy link
Copy Markdown
ContributorAuthor

Closing in favor of #30282

@mhencmhenc closed this Mar 24, 2023
@mhenc
mhenc deleted the ti_serialize branch August 30, 2023 18:48
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:APIAirflow's REST/HTTP APIarea:CLIarea:Schedulerincluding HA (high availability) schedulerarea:serializationarea:webserverWebserver related Issuesprovider:cncf-kubernetesKubernetes (k8s) provider related issues

Projects

None yet

Development

Successfully merging this pull request may close these issues.

7 participants

@mhenc@ashb@bolkedebruin@potiuk@uranusjr@kosteev@vincbeck
, '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

AIP-44 Support TaskInstance serialization/deserialization. - #29355

Closed
mhenc wants to merge 7 commits into
apache:mainfrom
mhenc:ti_serialize
Closed

AIP-44 Support TaskInstance serialization/deserialization.#29355
mhenc wants to merge 7 commits into
apache:mainfrom
mhenc:ti_serialize

Conversation

@mhenc

@mhencmhenc commented Feb 3, 2023

Copy link
Copy Markdown
Contributor

In AIP-44 we need to send the TaskInstance object to worker for execution. Currently this object type doesn't support serialization(only SimpleTaskInstance does).

This change is mainly about new way to construct TaskInstance from_dict - which requires changing the constructor into from_task and refactoring all usages.

Note that deserialized TaskInstance will not have ORM state (so will not be able to refresh state automatically from DB etc), but this should be fine as it will be used on client side of Internal API.

related: #29320

@boring-cyborgboring-cyborgBot added area:API Airflow's REST/HTTP API area:CLI area:core-operators provider:cncf-kubernetes Kubernetes (k8s) provider related issues area:Scheduler including HA (high availability) scheduler area:serialization area:webserver Webserver related Issues labels Feb 3, 2023
@mhenc
mhenc marked this pull request as ready for review February 3, 2023 14:16
@mhenc

mhenc commented Feb 3, 2023

Copy link
Copy Markdown
ContributorAuthor

cc: @potiuk@vincbeck

@mhencmhenc changed the title AIP-44 Support TaskInstance serialization.AIP-44 Support TaskInstance serialization/deserialization.Feb 3, 2023

@vincbeckvincbeck left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Nice!

Comment threadairflow/models/taskinstance.py Outdated

@kosteevkosteev left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

In order to achieve the goal of constructing TaskInstance objects from "task" and "dict"
could we have constructor of TaskInstance extended to support new parameter, e.g. task_dict like this:

def __init__(
self,
task: Operator,
execution_date: datetime | None = None,
run_id: str | None = None,
state: str | None = None,
map_index: int = -1,
task_dict: dict = None,
):
if task_dict is not None:
# init from task_dict
return
# old implementation

and then TaskInstance(task=t), TaskInstance(task_dict={...})

What do you think about this approach?

Comment threadairflow/models/taskinstance.py Outdated
@ashb

ashb commented Feb 8, 2023

Copy link
Copy Markdown
Member

@bolkedebruin could you take a look? I know you've recently overhauled all of the serialisation code

1 similar comment
@ashb

ashb commented Feb 8, 2023

Copy link
Copy Markdown
Member

@bolkedebruin could you take a look? I know you've recently overhauled all of the serialisation code

@bolkedebruinbolkedebruin left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Overall, great that this is happening! I would like to bring it more inline with the new serialization/deserialization code as I have been working on making it into 'one' (merging XCOM serialization and DAG serialization).

To seduce you: the new serialization code is an order of a magnitude faster than the old one (improvements over 50%) and can deal with versioning, while the old one cannot. The migration isn't fully completed yet (I am working on DAGs, but for obvious reasons that will take awhile), however that should not be in the way of making use of the new code (serde.py).

For your code this means moving away from to_dict/from_dict to the API of serialize() / deserialize(value, version) and setting __version__: ClassVar[int] = X in the TaskInstance class. This is pretty close to what you have now. serialization/deserialization should be called from serde and not directly. Alternatively you can add a serializer/deserializer to airflow.serialization.serializers.

I can help if needed.

Comment threadairflow/models/taskinstance.py Outdated
Comment threadairflow/models/taskinstance.py Outdated

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Here you are extending legacy code, i suggest using the more generic serialization code from serde

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

We rely on serialized_objects in InternalAPI in general
https://github.com/apache/airflow/blob/main/airflow/api_internal/internal_api_call.py#L107
Do you think we could switch to serde now? Is it compatible with what serialize_objects offer?

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Given that you are going to rely on a lot of serialization and the new serializer/deserializer is significantly faster and more future proof (versioning) I think everything that does currently not have a schema (that's everything except a DAG*) should switch.

I am willing to help out to ease the migration if required.

  • I am working on DAG serialization/deserialization but untangling how it is done now and to improve the structure is taking time especially with all the edge cases.

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

Great. I think we will migrate to the new serializer/deserializer in this case, but probably outside of this PR.
If you believe it would better to migrate first, then I can revert this change and get back to it when using new way.
WDYT?

@bolkedebruinbolkedebruinFeb 9, 2023

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Am I incorrect in thinking that migration is a 2 line change in internal_api_call and you are not relying on any of the other serialized_objects (basically DAG)? If so then I would say do not add technical dept and migrate now. This allows us to call serialized_objects as stale and soon to be deprecated.

Otherwise, keep it and and add it to the todo of AIP-44?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

yes, we need to make change there and server-side:
https://github.com/apache/airflow/blob/main/airflow/api_internal/endpoints/rpc_api_endpoint.py#L76

There are more methods with internal_api_call decorator. I did a quick check and I see that (beside primitives) we already need serialization to Dag,DagRun, BaseXCom, CallbackRequest (and probably more soon)

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.

Question. (and really sorry I have not looked at it earlier) and also follow up on #29513 as well.

Are we REALLY sure we need to serialize the whole TaskInstance (and other) objects to be passed via internal_api_call ?

For me this is an indication that we either have too narrow of a scope for an @internal_api call (generally speaking the whole internal_api_call should span the whole DB transaction. And since we are trying to pass an ORM object (TaskInstance, DAGRun etc.) it means that that object must have been retrieved before within a transaction. So it means that our internal_api_call should wrap the retrieval as well. Which might simply mean that we need to do some refactoring and add extract new methods (and then decorate them).

That's of course a general statement and approach and there might be cases that this require a bit deeper refactoring.

Which methods are affected @mhenc (besides the #29513 one) ? maybe we can look toghether and figure out approach for all of them ?

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.

For example in #29513 (comment) I proposed the solution that would avoid TaskInstance serialization altogether. I am reasonable convinced, that similar approach can be done for all ORM objects of ours and that we do not need to serialize any of them (in which case the whole PR might not be needed).

@mhencmhencFeb 21, 2023

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

But this works only for client -> Internal API side.

What about Internal API -> client, e.g. worker. We need to have TaskInstance object in worker to run the task, e.g.
https://github.com/apache/airflow/blob/main/airflow/cli/commands/task_command.py#L187

unless of course we are able to refactor it completely

@potiukpotiukFeb 21, 2023

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.

That's precisely what I am refactoring now (with my POC)

@bolkedebruin

bolkedebruin commented Feb 8, 2023

Copy link
Copy Markdown
Contributor

In order to achieve the goal of constructing TaskInstance objects from "task" and "dict" could we have constructor of TaskInstance extended to support new parameter, e.g. task_dict like this:

def init(
self,
task: Operator,
execution_date: datetime | None = None,
run_id: str | None = None,
state: str | None = None,
map_index: int = -1,
task_dict: dict = None,
):
if task_dict is not None:
# init from task_dict
return

# old implementation

and then TaskInstance(task=t), TaskInstance(task_dict={...})

What do you think about this approach?

I don't think that makes sense, a staticmethod or classmethod deserialize should cover this and used this way at other locations.

It could be as simple as this:

from airflow.serialization.serde import deserialize
ti = deserialize(serialized_ti)

@bolkedebruinbolkedebruin left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Added some additional comments

Comment threadairflow/models/taskinstance.py Outdated
Comment threadairflow/models/taskinstance.py Outdated
Comment threadairflow/models/taskinstance.py Outdated
@mhenc
mhenc requested review from bolkedebruin and kosteev and removed request for kosteevFebruary 9, 2023 13:38
@mhenc
mhenc requested a review from uranusjrFebruary 9, 2023 13:39

@bolkedebruinbolkedebruin left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Almost there :-)

Comment threadairflow/models/taskinstance.py Outdated
Comment threadairflow/models/taskinstance.py Outdated

@bolkedebruinbolkedebruinFeb 9, 2023

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Am I incorrect in thinking that migration is a 2 line change in internal_api_call and you are not relying on any of the other serialized_objects (basically DAG)? If so then I would say do not add technical dept and migrate now. This allows us to call serialized_objects as stale and soon to be deprecated.

Otherwise, keep it and and add it to the todo of AIP-44?

@mhenc

Copy link
Copy Markdown
ContributorAuthor

Closing in favor of #30282

@mhencmhenc closed this Mar 24, 2023
@mhenc
mhenc deleted the ti_serialize branch August 30, 2023 18:48
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:APIAirflow's REST/HTTP APIarea:CLIarea:Schedulerincluding HA (high availability) schedulerarea:serializationarea:webserverWebserver related Issuesprovider:cncf-kubernetesKubernetes (k8s) provider related issues

Projects

None yet

Development

Successfully merging this pull request may close these issues.

7 participants

@mhenc@ashb@bolkedebruin@potiuk@uranusjr@kosteev@vincbeck
, '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

AIP-44 Support TaskInstance serialization/deserialization. - #29355

Closed
mhenc wants to merge 7 commits into
apache:mainfrom
mhenc:ti_serialize
Closed

AIP-44 Support TaskInstance serialization/deserialization.#29355
mhenc wants to merge 7 commits into
apache:mainfrom
mhenc:ti_serialize

Conversation

@mhenc

@mhencmhenc commented Feb 3, 2023

Copy link
Copy Markdown
Contributor

In AIP-44 we need to send the TaskInstance object to worker for execution. Currently this object type doesn't support serialization(only SimpleTaskInstance does).

This change is mainly about new way to construct TaskInstance from_dict - which requires changing the constructor into from_task and refactoring all usages.

Note that deserialized TaskInstance will not have ORM state (so will not be able to refresh state automatically from DB etc), but this should be fine as it will be used on client side of Internal API.

related: #29320

@boring-cyborgboring-cyborgBot added area:API Airflow's REST/HTTP API area:CLI area:core-operators provider:cncf-kubernetes Kubernetes (k8s) provider related issues area:Scheduler including HA (high availability) scheduler area:serialization area:webserver Webserver related Issues labels Feb 3, 2023
@mhenc
mhenc marked this pull request as ready for review February 3, 2023 14:16
@mhenc

mhenc commented Feb 3, 2023

Copy link
Copy Markdown
ContributorAuthor

cc: @potiuk@vincbeck

@mhencmhenc changed the title AIP-44 Support TaskInstance serialization.AIP-44 Support TaskInstance serialization/deserialization.Feb 3, 2023

@vincbeckvincbeck left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Nice!

Comment threadairflow/models/taskinstance.py Outdated

@kosteevkosteev left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

In order to achieve the goal of constructing TaskInstance objects from "task" and "dict"
could we have constructor of TaskInstance extended to support new parameter, e.g. task_dict like this:

def __init__(
self,
task: Operator,
execution_date: datetime | None = None,
run_id: str | None = None,
state: str | None = None,
map_index: int = -1,
task_dict: dict = None,
):
if task_dict is not None:
# init from task_dict
return
# old implementation

and then TaskInstance(task=t), TaskInstance(task_dict={...})

What do you think about this approach?

Comment threadairflow/models/taskinstance.py Outdated
@ashb

ashb commented Feb 8, 2023

Copy link
Copy Markdown
Member

@bolkedebruin could you take a look? I know you've recently overhauled all of the serialisation code

1 similar comment
@ashb

ashb commented Feb 8, 2023

Copy link
Copy Markdown
Member

@bolkedebruin could you take a look? I know you've recently overhauled all of the serialisation code

@bolkedebruinbolkedebruin left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Overall, great that this is happening! I would like to bring it more inline with the new serialization/deserialization code as I have been working on making it into 'one' (merging XCOM serialization and DAG serialization).

To seduce you: the new serialization code is an order of a magnitude faster than the old one (improvements over 50%) and can deal with versioning, while the old one cannot. The migration isn't fully completed yet (I am working on DAGs, but for obvious reasons that will take awhile), however that should not be in the way of making use of the new code (serde.py).

For your code this means moving away from to_dict/from_dict to the API of serialize() / deserialize(value, version) and setting __version__: ClassVar[int] = X in the TaskInstance class. This is pretty close to what you have now. serialization/deserialization should be called from serde and not directly. Alternatively you can add a serializer/deserializer to airflow.serialization.serializers.

I can help if needed.

Comment threadairflow/models/taskinstance.py Outdated
Comment threadairflow/models/taskinstance.py Outdated

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Here you are extending legacy code, i suggest using the more generic serialization code from serde

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

We rely on serialized_objects in InternalAPI in general
https://github.com/apache/airflow/blob/main/airflow/api_internal/internal_api_call.py#L107
Do you think we could switch to serde now? Is it compatible with what serialize_objects offer?

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Given that you are going to rely on a lot of serialization and the new serializer/deserializer is significantly faster and more future proof (versioning) I think everything that does currently not have a schema (that's everything except a DAG*) should switch.

I am willing to help out to ease the migration if required.

  • I am working on DAG serialization/deserialization but untangling how it is done now and to improve the structure is taking time especially with all the edge cases.

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

Great. I think we will migrate to the new serializer/deserializer in this case, but probably outside of this PR.
If you believe it would better to migrate first, then I can revert this change and get back to it when using new way.
WDYT?

@bolkedebruinbolkedebruinFeb 9, 2023

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Am I incorrect in thinking that migration is a 2 line change in internal_api_call and you are not relying on any of the other serialized_objects (basically DAG)? If so then I would say do not add technical dept and migrate now. This allows us to call serialized_objects as stale and soon to be deprecated.

Otherwise, keep it and and add it to the todo of AIP-44?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

yes, we need to make change there and server-side:
https://github.com/apache/airflow/blob/main/airflow/api_internal/endpoints/rpc_api_endpoint.py#L76

There are more methods with internal_api_call decorator. I did a quick check and I see that (beside primitives) we already need serialization to Dag,DagRun, BaseXCom, CallbackRequest (and probably more soon)

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.

Question. (and really sorry I have not looked at it earlier) and also follow up on #29513 as well.

Are we REALLY sure we need to serialize the whole TaskInstance (and other) objects to be passed via internal_api_call ?

For me this is an indication that we either have too narrow of a scope for an @internal_api call (generally speaking the whole internal_api_call should span the whole DB transaction. And since we are trying to pass an ORM object (TaskInstance, DAGRun etc.) it means that that object must have been retrieved before within a transaction. So it means that our internal_api_call should wrap the retrieval as well. Which might simply mean that we need to do some refactoring and add extract new methods (and then decorate them).

That's of course a general statement and approach and there might be cases that this require a bit deeper refactoring.

Which methods are affected @mhenc (besides the #29513 one) ? maybe we can look toghether and figure out approach for all of them ?

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.

For example in #29513 (comment) I proposed the solution that would avoid TaskInstance serialization altogether. I am reasonable convinced, that similar approach can be done for all ORM objects of ours and that we do not need to serialize any of them (in which case the whole PR might not be needed).

@mhencmhencFeb 21, 2023

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

But this works only for client -> Internal API side.

What about Internal API -> client, e.g. worker. We need to have TaskInstance object in worker to run the task, e.g.
https://github.com/apache/airflow/blob/main/airflow/cli/commands/task_command.py#L187

unless of course we are able to refactor it completely

@potiukpotiukFeb 21, 2023

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.

That's precisely what I am refactoring now (with my POC)

@bolkedebruin

bolkedebruin commented Feb 8, 2023

Copy link
Copy Markdown
Contributor

In order to achieve the goal of constructing TaskInstance objects from "task" and "dict" could we have constructor of TaskInstance extended to support new parameter, e.g. task_dict like this:

def init(
self,
task: Operator,
execution_date: datetime | None = None,
run_id: str | None = None,
state: str | None = None,
map_index: int = -1,
task_dict: dict = None,
):
if task_dict is not None:
# init from task_dict
return

# old implementation

and then TaskInstance(task=t), TaskInstance(task_dict={...})

What do you think about this approach?

I don't think that makes sense, a staticmethod or classmethod deserialize should cover this and used this way at other locations.

It could be as simple as this:

from airflow.serialization.serde import deserialize
ti = deserialize(serialized_ti)

@bolkedebruinbolkedebruin left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Added some additional comments

Comment threadairflow/models/taskinstance.py Outdated
Comment threadairflow/models/taskinstance.py Outdated
Comment threadairflow/models/taskinstance.py Outdated
@mhenc
mhenc requested review from bolkedebruin and kosteev and removed request for kosteevFebruary 9, 2023 13:38
@mhenc
mhenc requested a review from uranusjrFebruary 9, 2023 13:39

@bolkedebruinbolkedebruin left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Almost there :-)

Comment threadairflow/models/taskinstance.py Outdated
Comment threadairflow/models/taskinstance.py Outdated

@bolkedebruinbolkedebruinFeb 9, 2023

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Am I incorrect in thinking that migration is a 2 line change in internal_api_call and you are not relying on any of the other serialized_objects (basically DAG)? If so then I would say do not add technical dept and migrate now. This allows us to call serialized_objects as stale and soon to be deprecated.

Otherwise, keep it and and add it to the todo of AIP-44?

@mhenc

Copy link
Copy Markdown
ContributorAuthor

Closing in favor of #30282

@mhencmhenc closed this Mar 24, 2023
@mhenc
mhenc deleted the ti_serialize branch August 30, 2023 18:48
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:APIAirflow's REST/HTTP APIarea:CLIarea:Schedulerincluding HA (high availability) schedulerarea:serializationarea:webserverWebserver related Issuesprovider:cncf-kubernetesKubernetes (k8s) provider related issues

Projects

None yet

Development

Successfully merging this pull request may close these issues.

7 participants

@mhenc@ashb@bolkedebruin@potiuk@uranusjr@kosteev@vincbeck
, '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

AIP-44 Support TaskInstance serialization/deserialization. - #29355

Closed
mhenc wants to merge 7 commits into
apache:mainfrom
mhenc:ti_serialize
Closed

AIP-44 Support TaskInstance serialization/deserialization.#29355
mhenc wants to merge 7 commits into
apache:mainfrom
mhenc:ti_serialize

Conversation

@mhenc

@mhencmhenc commented Feb 3, 2023

Copy link
Copy Markdown
Contributor

In AIP-44 we need to send the TaskInstance object to worker for execution. Currently this object type doesn't support serialization(only SimpleTaskInstance does).

This change is mainly about new way to construct TaskInstance from_dict - which requires changing the constructor into from_task and refactoring all usages.

Note that deserialized TaskInstance will not have ORM state (so will not be able to refresh state automatically from DB etc), but this should be fine as it will be used on client side of Internal API.

related: #29320

@boring-cyborgboring-cyborgBot added area:API Airflow's REST/HTTP API area:CLI area:core-operators provider:cncf-kubernetes Kubernetes (k8s) provider related issues area:Scheduler including HA (high availability) scheduler area:serialization area:webserver Webserver related Issues labels Feb 3, 2023
@mhenc
mhenc marked this pull request as ready for review February 3, 2023 14:16
@mhenc

mhenc commented Feb 3, 2023

Copy link
Copy Markdown
ContributorAuthor

cc: @potiuk@vincbeck

@mhencmhenc changed the title AIP-44 Support TaskInstance serialization.AIP-44 Support TaskInstance serialization/deserialization.Feb 3, 2023

@vincbeckvincbeck left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Nice!

Comment threadairflow/models/taskinstance.py Outdated

@kosteevkosteev left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

In order to achieve the goal of constructing TaskInstance objects from "task" and "dict"
could we have constructor of TaskInstance extended to support new parameter, e.g. task_dict like this:

def __init__(
self,
task: Operator,
execution_date: datetime | None = None,
run_id: str | None = None,
state: str | None = None,
map_index: int = -1,
task_dict: dict = None,
):
if task_dict is not None:
# init from task_dict
return
# old implementation

and then TaskInstance(task=t), TaskInstance(task_dict={...})

What do you think about this approach?

Comment threadairflow/models/taskinstance.py Outdated
@ashb

ashb commented Feb 8, 2023

Copy link
Copy Markdown
Member

@bolkedebruin could you take a look? I know you've recently overhauled all of the serialisation code

1 similar comment
@ashb

ashb commented Feb 8, 2023

Copy link
Copy Markdown
Member

@bolkedebruin could you take a look? I know you've recently overhauled all of the serialisation code

@bolkedebruinbolkedebruin left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Overall, great that this is happening! I would like to bring it more inline with the new serialization/deserialization code as I have been working on making it into 'one' (merging XCOM serialization and DAG serialization).

To seduce you: the new serialization code is an order of a magnitude faster than the old one (improvements over 50%) and can deal with versioning, while the old one cannot. The migration isn't fully completed yet (I am working on DAGs, but for obvious reasons that will take awhile), however that should not be in the way of making use of the new code (serde.py).

For your code this means moving away from to_dict/from_dict to the API of serialize() / deserialize(value, version) and setting __version__: ClassVar[int] = X in the TaskInstance class. This is pretty close to what you have now. serialization/deserialization should be called from serde and not directly. Alternatively you can add a serializer/deserializer to airflow.serialization.serializers.

I can help if needed.

Comment threadairflow/models/taskinstance.py Outdated
Comment threadairflow/models/taskinstance.py Outdated

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Here you are extending legacy code, i suggest using the more generic serialization code from serde

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

We rely on serialized_objects in InternalAPI in general
https://github.com/apache/airflow/blob/main/airflow/api_internal/internal_api_call.py#L107
Do you think we could switch to serde now? Is it compatible with what serialize_objects offer?

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Given that you are going to rely on a lot of serialization and the new serializer/deserializer is significantly faster and more future proof (versioning) I think everything that does currently not have a schema (that's everything except a DAG*) should switch.

I am willing to help out to ease the migration if required.

  • I am working on DAG serialization/deserialization but untangling how it is done now and to improve the structure is taking time especially with all the edge cases.

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

Great. I think we will migrate to the new serializer/deserializer in this case, but probably outside of this PR.
If you believe it would better to migrate first, then I can revert this change and get back to it when using new way.
WDYT?

@bolkedebruinbolkedebruinFeb 9, 2023

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Am I incorrect in thinking that migration is a 2 line change in internal_api_call and you are not relying on any of the other serialized_objects (basically DAG)? If so then I would say do not add technical dept and migrate now. This allows us to call serialized_objects as stale and soon to be deprecated.

Otherwise, keep it and and add it to the todo of AIP-44?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

yes, we need to make change there and server-side:
https://github.com/apache/airflow/blob/main/airflow/api_internal/endpoints/rpc_api_endpoint.py#L76

There are more methods with internal_api_call decorator. I did a quick check and I see that (beside primitives) we already need serialization to Dag,DagRun, BaseXCom, CallbackRequest (and probably more soon)

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.

Question. (and really sorry I have not looked at it earlier) and also follow up on #29513 as well.

Are we REALLY sure we need to serialize the whole TaskInstance (and other) objects to be passed via internal_api_call ?

For me this is an indication that we either have too narrow of a scope for an @internal_api call (generally speaking the whole internal_api_call should span the whole DB transaction. And since we are trying to pass an ORM object (TaskInstance, DAGRun etc.) it means that that object must have been retrieved before within a transaction. So it means that our internal_api_call should wrap the retrieval as well. Which might simply mean that we need to do some refactoring and add extract new methods (and then decorate them).

That's of course a general statement and approach and there might be cases that this require a bit deeper refactoring.

Which methods are affected @mhenc (besides the #29513 one) ? maybe we can look toghether and figure out approach for all of them ?

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.

For example in #29513 (comment) I proposed the solution that would avoid TaskInstance serialization altogether. I am reasonable convinced, that similar approach can be done for all ORM objects of ours and that we do not need to serialize any of them (in which case the whole PR might not be needed).

@mhencmhencFeb 21, 2023

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

But this works only for client -> Internal API side.

What about Internal API -> client, e.g. worker. We need to have TaskInstance object in worker to run the task, e.g.
https://github.com/apache/airflow/blob/main/airflow/cli/commands/task_command.py#L187

unless of course we are able to refactor it completely

@potiukpotiukFeb 21, 2023

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.

That's precisely what I am refactoring now (with my POC)

@bolkedebruin

bolkedebruin commented Feb 8, 2023

Copy link
Copy Markdown
Contributor

In order to achieve the goal of constructing TaskInstance objects from "task" and "dict" could we have constructor of TaskInstance extended to support new parameter, e.g. task_dict like this:

def init(
self,
task: Operator,
execution_date: datetime | None = None,
run_id: str | None = None,
state: str | None = None,
map_index: int = -1,
task_dict: dict = None,
):
if task_dict is not None:
# init from task_dict
return

# old implementation

and then TaskInstance(task=t), TaskInstance(task_dict={...})

What do you think about this approach?

I don't think that makes sense, a staticmethod or classmethod deserialize should cover this and used this way at other locations.

It could be as simple as this:

from airflow.serialization.serde import deserialize
ti = deserialize(serialized_ti)

@bolkedebruinbolkedebruin left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Added some additional comments

Comment threadairflow/models/taskinstance.py Outdated
Comment threadairflow/models/taskinstance.py Outdated
Comment threadairflow/models/taskinstance.py Outdated
@mhenc
mhenc requested review from bolkedebruin and kosteev and removed request for kosteevFebruary 9, 2023 13:38
@mhenc
mhenc requested a review from uranusjrFebruary 9, 2023 13:39

@bolkedebruinbolkedebruin left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Almost there :-)

Comment threadairflow/models/taskinstance.py Outdated
Comment threadairflow/models/taskinstance.py Outdated

@bolkedebruinbolkedebruinFeb 9, 2023

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Am I incorrect in thinking that migration is a 2 line change in internal_api_call and you are not relying on any of the other serialized_objects (basically DAG)? If so then I would say do not add technical dept and migrate now. This allows us to call serialized_objects as stale and soon to be deprecated.

Otherwise, keep it and and add it to the todo of AIP-44?

@mhenc

Copy link
Copy Markdown
ContributorAuthor

Closing in favor of #30282

@mhencmhenc closed this Mar 24, 2023
@mhenc
mhenc deleted the ti_serialize branch August 30, 2023 18:48
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:APIAirflow's REST/HTTP APIarea:CLIarea:Schedulerincluding HA (high availability) schedulerarea:serializationarea:webserverWebserver related Issuesprovider:cncf-kubernetesKubernetes (k8s) provider related issues

Projects

None yet

Development

Successfully merging this pull request may close these issues.

7 participants

@mhenc@ashb@bolkedebruin@potiuk@uranusjr@kosteev@vincbeck
, '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

AIP-44 Support TaskInstance serialization/deserialization. - #29355

Closed
mhenc wants to merge 7 commits into
apache:mainfrom
mhenc:ti_serialize
Closed

AIP-44 Support TaskInstance serialization/deserialization.#29355
mhenc wants to merge 7 commits into
apache:mainfrom
mhenc:ti_serialize

Conversation

@mhenc

@mhencmhenc commented Feb 3, 2023

Copy link
Copy Markdown
Contributor

In AIP-44 we need to send the TaskInstance object to worker for execution. Currently this object type doesn't support serialization(only SimpleTaskInstance does).

This change is mainly about new way to construct TaskInstance from_dict - which requires changing the constructor into from_task and refactoring all usages.

Note that deserialized TaskInstance will not have ORM state (so will not be able to refresh state automatically from DB etc), but this should be fine as it will be used on client side of Internal API.

related: #29320

@boring-cyborgboring-cyborgBot added area:API Airflow's REST/HTTP API area:CLI area:core-operators provider:cncf-kubernetes Kubernetes (k8s) provider related issues area:Scheduler including HA (high availability) scheduler area:serialization area:webserver Webserver related Issues labels Feb 3, 2023
@mhenc
mhenc marked this pull request as ready for review February 3, 2023 14:16
@mhenc

mhenc commented Feb 3, 2023

Copy link
Copy Markdown
ContributorAuthor

cc: @potiuk@vincbeck

@mhencmhenc changed the title AIP-44 Support TaskInstance serialization.AIP-44 Support TaskInstance serialization/deserialization.Feb 3, 2023

@vincbeckvincbeck left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Nice!

Comment threadairflow/models/taskinstance.py Outdated

@kosteevkosteev left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

In order to achieve the goal of constructing TaskInstance objects from "task" and "dict"
could we have constructor of TaskInstance extended to support new parameter, e.g. task_dict like this:

def __init__(
self,
task: Operator,
execution_date: datetime | None = None,
run_id: str | None = None,
state: str | None = None,
map_index: int = -1,
task_dict: dict = None,
):
if task_dict is not None:
# init from task_dict
return
# old implementation

and then TaskInstance(task=t), TaskInstance(task_dict={...})

What do you think about this approach?

Comment threadairflow/models/taskinstance.py Outdated
@ashb

ashb commented Feb 8, 2023

Copy link
Copy Markdown
Member

@bolkedebruin could you take a look? I know you've recently overhauled all of the serialisation code

1 similar comment
@ashb

ashb commented Feb 8, 2023

Copy link
Copy Markdown
Member

@bolkedebruin could you take a look? I know you've recently overhauled all of the serialisation code

@bolkedebruinbolkedebruin left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Overall, great that this is happening! I would like to bring it more inline with the new serialization/deserialization code as I have been working on making it into 'one' (merging XCOM serialization and DAG serialization).

To seduce you: the new serialization code is an order of a magnitude faster than the old one (improvements over 50%) and can deal with versioning, while the old one cannot. The migration isn't fully completed yet (I am working on DAGs, but for obvious reasons that will take awhile), however that should not be in the way of making use of the new code (serde.py).

For your code this means moving away from to_dict/from_dict to the API of serialize() / deserialize(value, version) and setting __version__: ClassVar[int] = X in the TaskInstance class. This is pretty close to what you have now. serialization/deserialization should be called from serde and not directly. Alternatively you can add a serializer/deserializer to airflow.serialization.serializers.

I can help if needed.

Comment threadairflow/models/taskinstance.py Outdated
Comment threadairflow/models/taskinstance.py Outdated

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Here you are extending legacy code, i suggest using the more generic serialization code from serde

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

We rely on serialized_objects in InternalAPI in general
https://github.com/apache/airflow/blob/main/airflow/api_internal/internal_api_call.py#L107
Do you think we could switch to serde now? Is it compatible with what serialize_objects offer?

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Given that you are going to rely on a lot of serialization and the new serializer/deserializer is significantly faster and more future proof (versioning) I think everything that does currently not have a schema (that's everything except a DAG*) should switch.

I am willing to help out to ease the migration if required.

  • I am working on DAG serialization/deserialization but untangling how it is done now and to improve the structure is taking time especially with all the edge cases.

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

Great. I think we will migrate to the new serializer/deserializer in this case, but probably outside of this PR.
If you believe it would better to migrate first, then I can revert this change and get back to it when using new way.
WDYT?

@bolkedebruinbolkedebruinFeb 9, 2023

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Am I incorrect in thinking that migration is a 2 line change in internal_api_call and you are not relying on any of the other serialized_objects (basically DAG)? If so then I would say do not add technical dept and migrate now. This allows us to call serialized_objects as stale and soon to be deprecated.

Otherwise, keep it and and add it to the todo of AIP-44?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

yes, we need to make change there and server-side:
https://github.com/apache/airflow/blob/main/airflow/api_internal/endpoints/rpc_api_endpoint.py#L76

There are more methods with internal_api_call decorator. I did a quick check and I see that (beside primitives) we already need serialization to Dag,DagRun, BaseXCom, CallbackRequest (and probably more soon)

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.

Question. (and really sorry I have not looked at it earlier) and also follow up on #29513 as well.

Are we REALLY sure we need to serialize the whole TaskInstance (and other) objects to be passed via internal_api_call ?

For me this is an indication that we either have too narrow of a scope for an @internal_api call (generally speaking the whole internal_api_call should span the whole DB transaction. And since we are trying to pass an ORM object (TaskInstance, DAGRun etc.) it means that that object must have been retrieved before within a transaction. So it means that our internal_api_call should wrap the retrieval as well. Which might simply mean that we need to do some refactoring and add extract new methods (and then decorate them).

That's of course a general statement and approach and there might be cases that this require a bit deeper refactoring.

Which methods are affected @mhenc (besides the #29513 one) ? maybe we can look toghether and figure out approach for all of them ?

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.

For example in #29513 (comment) I proposed the solution that would avoid TaskInstance serialization altogether. I am reasonable convinced, that similar approach can be done for all ORM objects of ours and that we do not need to serialize any of them (in which case the whole PR might not be needed).

@mhencmhencFeb 21, 2023

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

But this works only for client -> Internal API side.

What about Internal API -> client, e.g. worker. We need to have TaskInstance object in worker to run the task, e.g.
https://github.com/apache/airflow/blob/main/airflow/cli/commands/task_command.py#L187

unless of course we are able to refactor it completely

@potiukpotiukFeb 21, 2023

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.

That's precisely what I am refactoring now (with my POC)

@bolkedebruin

bolkedebruin commented Feb 8, 2023

Copy link
Copy Markdown
Contributor

In order to achieve the goal of constructing TaskInstance objects from "task" and "dict" could we have constructor of TaskInstance extended to support new parameter, e.g. task_dict like this:

def init(
self,
task: Operator,
execution_date: datetime | None = None,
run_id: str | None = None,
state: str | None = None,
map_index: int = -1,
task_dict: dict = None,
):
if task_dict is not None:
# init from task_dict
return

# old implementation

and then TaskInstance(task=t), TaskInstance(task_dict={...})

What do you think about this approach?

I don't think that makes sense, a staticmethod or classmethod deserialize should cover this and used this way at other locations.

It could be as simple as this:

from airflow.serialization.serde import deserialize
ti = deserialize(serialized_ti)

@bolkedebruinbolkedebruin left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Added some additional comments

Comment threadairflow/models/taskinstance.py Outdated
Comment threadairflow/models/taskinstance.py Outdated
Comment threadairflow/models/taskinstance.py Outdated
@mhenc
mhenc requested review from bolkedebruin and kosteev and removed request for kosteevFebruary 9, 2023 13:38
@mhenc
mhenc requested a review from uranusjrFebruary 9, 2023 13:39

@bolkedebruinbolkedebruin left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Almost there :-)

Comment threadairflow/models/taskinstance.py Outdated
Comment threadairflow/models/taskinstance.py Outdated

@bolkedebruinbolkedebruinFeb 9, 2023

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Am I incorrect in thinking that migration is a 2 line change in internal_api_call and you are not relying on any of the other serialized_objects (basically DAG)? If so then I would say do not add technical dept and migrate now. This allows us to call serialized_objects as stale and soon to be deprecated.

Otherwise, keep it and and add it to the todo of AIP-44?

@mhenc

Copy link
Copy Markdown
ContributorAuthor

Closing in favor of #30282

@mhencmhenc closed this Mar 24, 2023
@mhenc
mhenc deleted the ti_serialize branch August 30, 2023 18:48
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:APIAirflow's REST/HTTP APIarea:CLIarea:Schedulerincluding HA (high availability) schedulerarea:serializationarea:webserverWebserver related Issuesprovider:cncf-kubernetesKubernetes (k8s) provider related issues

Projects

None yet

Development

Successfully merging this pull request may close these issues.

7 participants

@mhenc@ashb@bolkedebruin@potiuk@uranusjr@kosteev@vincbeck
, '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

AIP-44 Support TaskInstance serialization/deserialization. - #29355

Closed
mhenc wants to merge 7 commits into
apache:mainfrom
mhenc:ti_serialize
Closed

AIP-44 Support TaskInstance serialization/deserialization.#29355
mhenc wants to merge 7 commits into
apache:mainfrom
mhenc:ti_serialize

Conversation

@mhenc

@mhencmhenc commented Feb 3, 2023

Copy link
Copy Markdown
Contributor

In AIP-44 we need to send the TaskInstance object to worker for execution. Currently this object type doesn't support serialization(only SimpleTaskInstance does).

This change is mainly about new way to construct TaskInstance from_dict - which requires changing the constructor into from_task and refactoring all usages.

Note that deserialized TaskInstance will not have ORM state (so will not be able to refresh state automatically from DB etc), but this should be fine as it will be used on client side of Internal API.

related: #29320

@boring-cyborgboring-cyborgBot added area:API Airflow's REST/HTTP API area:CLI area:core-operators provider:cncf-kubernetes Kubernetes (k8s) provider related issues area:Scheduler including HA (high availability) scheduler area:serialization area:webserver Webserver related Issues labels Feb 3, 2023
@mhenc
mhenc marked this pull request as ready for review February 3, 2023 14:16
@mhenc

mhenc commented Feb 3, 2023

Copy link
Copy Markdown
ContributorAuthor

cc: @potiuk@vincbeck

@mhencmhenc changed the title AIP-44 Support TaskInstance serialization.AIP-44 Support TaskInstance serialization/deserialization.Feb 3, 2023

@vincbeckvincbeck left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Nice!

Comment threadairflow/models/taskinstance.py Outdated

@kosteevkosteev left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

In order to achieve the goal of constructing TaskInstance objects from "task" and "dict"
could we have constructor of TaskInstance extended to support new parameter, e.g. task_dict like this:

def __init__(
self,
task: Operator,
execution_date: datetime | None = None,
run_id: str | None = None,
state: str | None = None,
map_index: int = -1,
task_dict: dict = None,
):
if task_dict is not None:
# init from task_dict
return
# old implementation

and then TaskInstance(task=t), TaskInstance(task_dict={...})

What do you think about this approach?

Comment threadairflow/models/taskinstance.py Outdated
@ashb

ashb commented Feb 8, 2023

Copy link
Copy Markdown
Member

@bolkedebruin could you take a look? I know you've recently overhauled all of the serialisation code

1 similar comment
@ashb

ashb commented Feb 8, 2023

Copy link
Copy Markdown
Member

@bolkedebruin could you take a look? I know you've recently overhauled all of the serialisation code

@bolkedebruinbolkedebruin left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Overall, great that this is happening! I would like to bring it more inline with the new serialization/deserialization code as I have been working on making it into 'one' (merging XCOM serialization and DAG serialization).

To seduce you: the new serialization code is an order of a magnitude faster than the old one (improvements over 50%) and can deal with versioning, while the old one cannot. The migration isn't fully completed yet (I am working on DAGs, but for obvious reasons that will take awhile), however that should not be in the way of making use of the new code (serde.py).

For your code this means moving away from to_dict/from_dict to the API of serialize() / deserialize(value, version) and setting __version__: ClassVar[int] = X in the TaskInstance class. This is pretty close to what you have now. serialization/deserialization should be called from serde and not directly. Alternatively you can add a serializer/deserializer to airflow.serialization.serializers.

I can help if needed.

Comment threadairflow/models/taskinstance.py Outdated
Comment threadairflow/models/taskinstance.py Outdated

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Here you are extending legacy code, i suggest using the more generic serialization code from serde

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

We rely on serialized_objects in InternalAPI in general
https://github.com/apache/airflow/blob/main/airflow/api_internal/internal_api_call.py#L107
Do you think we could switch to serde now? Is it compatible with what serialize_objects offer?

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Given that you are going to rely on a lot of serialization and the new serializer/deserializer is significantly faster and more future proof (versioning) I think everything that does currently not have a schema (that's everything except a DAG*) should switch.

I am willing to help out to ease the migration if required.

  • I am working on DAG serialization/deserialization but untangling how it is done now and to improve the structure is taking time especially with all the edge cases.

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

Great. I think we will migrate to the new serializer/deserializer in this case, but probably outside of this PR.
If you believe it would better to migrate first, then I can revert this change and get back to it when using new way.
WDYT?

@bolkedebruinbolkedebruinFeb 9, 2023

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Am I incorrect in thinking that migration is a 2 line change in internal_api_call and you are not relying on any of the other serialized_objects (basically DAG)? If so then I would say do not add technical dept and migrate now. This allows us to call serialized_objects as stale and soon to be deprecated.

Otherwise, keep it and and add it to the todo of AIP-44?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

yes, we need to make change there and server-side:
https://github.com/apache/airflow/blob/main/airflow/api_internal/endpoints/rpc_api_endpoint.py#L76

There are more methods with internal_api_call decorator. I did a quick check and I see that (beside primitives) we already need serialization to Dag,DagRun, BaseXCom, CallbackRequest (and probably more soon)

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.

Question. (and really sorry I have not looked at it earlier) and also follow up on #29513 as well.

Are we REALLY sure we need to serialize the whole TaskInstance (and other) objects to be passed via internal_api_call ?

For me this is an indication that we either have too narrow of a scope for an @internal_api call (generally speaking the whole internal_api_call should span the whole DB transaction. And since we are trying to pass an ORM object (TaskInstance, DAGRun etc.) it means that that object must have been retrieved before within a transaction. So it means that our internal_api_call should wrap the retrieval as well. Which might simply mean that we need to do some refactoring and add extract new methods (and then decorate them).

That's of course a general statement and approach and there might be cases that this require a bit deeper refactoring.

Which methods are affected @mhenc (besides the #29513 one) ? maybe we can look toghether and figure out approach for all of them ?

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.

For example in #29513 (comment) I proposed the solution that would avoid TaskInstance serialization altogether. I am reasonable convinced, that similar approach can be done for all ORM objects of ours and that we do not need to serialize any of them (in which case the whole PR might not be needed).

@mhencmhencFeb 21, 2023

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

But this works only for client -> Internal API side.

What about Internal API -> client, e.g. worker. We need to have TaskInstance object in worker to run the task, e.g.
https://github.com/apache/airflow/blob/main/airflow/cli/commands/task_command.py#L187

unless of course we are able to refactor it completely

@potiukpotiukFeb 21, 2023

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.

That's precisely what I am refactoring now (with my POC)

@bolkedebruin

bolkedebruin commented Feb 8, 2023

Copy link
Copy Markdown
Contributor

In order to achieve the goal of constructing TaskInstance objects from "task" and "dict" could we have constructor of TaskInstance extended to support new parameter, e.g. task_dict like this:

def init(
self,
task: Operator,
execution_date: datetime | None = None,
run_id: str | None = None,
state: str | None = None,
map_index: int = -1,
task_dict: dict = None,
):
if task_dict is not None:
# init from task_dict
return

# old implementation

and then TaskInstance(task=t), TaskInstance(task_dict={...})

What do you think about this approach?

I don't think that makes sense, a staticmethod or classmethod deserialize should cover this and used this way at other locations.

It could be as simple as this:

from airflow.serialization.serde import deserialize
ti = deserialize(serialized_ti)

@bolkedebruinbolkedebruin left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Added some additional comments

Comment threadairflow/models/taskinstance.py Outdated
Comment threadairflow/models/taskinstance.py Outdated
Comment threadairflow/models/taskinstance.py Outdated
@mhenc
mhenc requested review from bolkedebruin and kosteev and removed request for kosteevFebruary 9, 2023 13:38
@mhenc
mhenc requested a review from uranusjrFebruary 9, 2023 13:39

@bolkedebruinbolkedebruin left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Almost there :-)

Comment threadairflow/models/taskinstance.py Outdated
Comment threadairflow/models/taskinstance.py Outdated

@bolkedebruinbolkedebruinFeb 9, 2023

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Am I incorrect in thinking that migration is a 2 line change in internal_api_call and you are not relying on any of the other serialized_objects (basically DAG)? If so then I would say do not add technical dept and migrate now. This allows us to call serialized_objects as stale and soon to be deprecated.

Otherwise, keep it and and add it to the todo of AIP-44?

@mhenc

Copy link
Copy Markdown
ContributorAuthor

Closing in favor of #30282

@mhencmhenc closed this Mar 24, 2023
@mhenc
mhenc deleted the ti_serialize branch August 30, 2023 18:48
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:APIAirflow's REST/HTTP APIarea:CLIarea:Schedulerincluding HA (high availability) schedulerarea:serializationarea:webserverWebserver related Issuesprovider:cncf-kubernetesKubernetes (k8s) provider related issues

Projects

None yet

Development

Successfully merging this pull request may close these issues.

7 participants

@mhenc@ashb@bolkedebruin@potiuk@uranusjr@kosteev@vincbeck
, '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

AIP-44 Support TaskInstance serialization/deserialization. - #29355

Closed
mhenc wants to merge 7 commits into
apache:mainfrom
mhenc:ti_serialize
Closed

AIP-44 Support TaskInstance serialization/deserialization.#29355
mhenc wants to merge 7 commits into
apache:mainfrom
mhenc:ti_serialize

Conversation

@mhenc

@mhencmhenc commented Feb 3, 2023

Copy link
Copy Markdown
Contributor

In AIP-44 we need to send the TaskInstance object to worker for execution. Currently this object type doesn't support serialization(only SimpleTaskInstance does).

This change is mainly about new way to construct TaskInstance from_dict - which requires changing the constructor into from_task and refactoring all usages.

Note that deserialized TaskInstance will not have ORM state (so will not be able to refresh state automatically from DB etc), but this should be fine as it will be used on client side of Internal API.

related: #29320

@boring-cyborgboring-cyborgBot added area:API Airflow's REST/HTTP API area:CLI area:core-operators provider:cncf-kubernetes Kubernetes (k8s) provider related issues area:Scheduler including HA (high availability) scheduler area:serialization area:webserver Webserver related Issues labels Feb 3, 2023
@mhenc
mhenc marked this pull request as ready for review February 3, 2023 14:16
@mhenc

mhenc commented Feb 3, 2023

Copy link
Copy Markdown
ContributorAuthor

cc: @potiuk@vincbeck

@mhencmhenc changed the title AIP-44 Support TaskInstance serialization.AIP-44 Support TaskInstance serialization/deserialization.Feb 3, 2023

@vincbeckvincbeck left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Nice!

Comment threadairflow/models/taskinstance.py Outdated

@kosteevkosteev left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

In order to achieve the goal of constructing TaskInstance objects from "task" and "dict"
could we have constructor of TaskInstance extended to support new parameter, e.g. task_dict like this:

def __init__(
self,
task: Operator,
execution_date: datetime | None = None,
run_id: str | None = None,
state: str | None = None,
map_index: int = -1,
task_dict: dict = None,
):
if task_dict is not None:
# init from task_dict
return
# old implementation

and then TaskInstance(task=t), TaskInstance(task_dict={...})

What do you think about this approach?

Comment threadairflow/models/taskinstance.py Outdated
@ashb

ashb commented Feb 8, 2023

Copy link
Copy Markdown
Member

@bolkedebruin could you take a look? I know you've recently overhauled all of the serialisation code

1 similar comment
@ashb

ashb commented Feb 8, 2023

Copy link
Copy Markdown
Member

@bolkedebruin could you take a look? I know you've recently overhauled all of the serialisation code

@bolkedebruinbolkedebruin left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Overall, great that this is happening! I would like to bring it more inline with the new serialization/deserialization code as I have been working on making it into 'one' (merging XCOM serialization and DAG serialization).

To seduce you: the new serialization code is an order of a magnitude faster than the old one (improvements over 50%) and can deal with versioning, while the old one cannot. The migration isn't fully completed yet (I am working on DAGs, but for obvious reasons that will take awhile), however that should not be in the way of making use of the new code (serde.py).

For your code this means moving away from to_dict/from_dict to the API of serialize() / deserialize(value, version) and setting __version__: ClassVar[int] = X in the TaskInstance class. This is pretty close to what you have now. serialization/deserialization should be called from serde and not directly. Alternatively you can add a serializer/deserializer to airflow.serialization.serializers.

I can help if needed.

Comment threadairflow/models/taskinstance.py Outdated
Comment threadairflow/models/taskinstance.py Outdated

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Here you are extending legacy code, i suggest using the more generic serialization code from serde

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

We rely on serialized_objects in InternalAPI in general
https://github.com/apache/airflow/blob/main/airflow/api_internal/internal_api_call.py#L107
Do you think we could switch to serde now? Is it compatible with what serialize_objects offer?

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Given that you are going to rely on a lot of serialization and the new serializer/deserializer is significantly faster and more future proof (versioning) I think everything that does currently not have a schema (that's everything except a DAG*) should switch.

I am willing to help out to ease the migration if required.

  • I am working on DAG serialization/deserialization but untangling how it is done now and to improve the structure is taking time especially with all the edge cases.

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

Great. I think we will migrate to the new serializer/deserializer in this case, but probably outside of this PR.
If you believe it would better to migrate first, then I can revert this change and get back to it when using new way.
WDYT?

@bolkedebruinbolkedebruinFeb 9, 2023

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Am I incorrect in thinking that migration is a 2 line change in internal_api_call and you are not relying on any of the other serialized_objects (basically DAG)? If so then I would say do not add technical dept and migrate now. This allows us to call serialized_objects as stale and soon to be deprecated.

Otherwise, keep it and and add it to the todo of AIP-44?

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

yes, we need to make change there and server-side:
https://github.com/apache/airflow/blob/main/airflow/api_internal/endpoints/rpc_api_endpoint.py#L76

There are more methods with internal_api_call decorator. I did a quick check and I see that (beside primitives) we already need serialization to Dag,DagRun, BaseXCom, CallbackRequest (and probably more soon)

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.

Question. (and really sorry I have not looked at it earlier) and also follow up on #29513 as well.

Are we REALLY sure we need to serialize the whole TaskInstance (and other) objects to be passed via internal_api_call ?

For me this is an indication that we either have too narrow of a scope for an @internal_api call (generally speaking the whole internal_api_call should span the whole DB transaction. And since we are trying to pass an ORM object (TaskInstance, DAGRun etc.) it means that that object must have been retrieved before within a transaction. So it means that our internal_api_call should wrap the retrieval as well. Which might simply mean that we need to do some refactoring and add extract new methods (and then decorate them).

That's of course a general statement and approach and there might be cases that this require a bit deeper refactoring.

Which methods are affected @mhenc (besides the #29513 one) ? maybe we can look toghether and figure out approach for all of them ?

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.

For example in #29513 (comment) I proposed the solution that would avoid TaskInstance serialization altogether. I am reasonable convinced, that similar approach can be done for all ORM objects of ours and that we do not need to serialize any of them (in which case the whole PR might not be needed).

@mhencmhencFeb 21, 2023

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

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

But this works only for client -> Internal API side.

What about Internal API -> client, e.g. worker. We need to have TaskInstance object in worker to run the task, e.g.
https://github.com/apache/airflow/blob/main/airflow/cli/commands/task_command.py#L187

unless of course we are able to refactor it completely

@potiukpotiukFeb 21, 2023

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.

That's precisely what I am refactoring now (with my POC)

@bolkedebruin

bolkedebruin commented Feb 8, 2023

Copy link
Copy Markdown
Contributor

In order to achieve the goal of constructing TaskInstance objects from "task" and "dict" could we have constructor of TaskInstance extended to support new parameter, e.g. task_dict like this:

def init(
self,
task: Operator,
execution_date: datetime | None = None,
run_id: str | None = None,
state: str | None = None,
map_index: int = -1,
task_dict: dict = None,
):
if task_dict is not None:
# init from task_dict
return

# old implementation

and then TaskInstance(task=t), TaskInstance(task_dict={...})

What do you think about this approach?

I don't think that makes sense, a staticmethod or classmethod deserialize should cover this and used this way at other locations.

It could be as simple as this:

from airflow.serialization.serde import deserialize
ti = deserialize(serialized_ti)

@bolkedebruinbolkedebruin left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Added some additional comments

Comment threadairflow/models/taskinstance.py Outdated
Comment threadairflow/models/taskinstance.py Outdated
Comment threadairflow/models/taskinstance.py Outdated
@mhenc
mhenc requested review from bolkedebruin and kosteev and removed request for kosteevFebruary 9, 2023 13:38
@mhenc
mhenc requested a review from uranusjrFebruary 9, 2023 13:39

@bolkedebruinbolkedebruin left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Almost there :-)

Comment threadairflow/models/taskinstance.py Outdated
Comment threadairflow/models/taskinstance.py Outdated

@bolkedebruinbolkedebruinFeb 9, 2023

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Am I incorrect in thinking that migration is a 2 line change in internal_api_call and you are not relying on any of the other serialized_objects (basically DAG)? If so then I would say do not add technical dept and migrate now. This allows us to call serialized_objects as stale and soon to be deprecated.

Otherwise, keep it and and add it to the todo of AIP-44?

@mhenc

Copy link
Copy Markdown
ContributorAuthor

Closing in favor of #30282

@mhencmhenc closed this Mar 24, 2023
@mhenc
mhenc deleted the ti_serialize branch August 30, 2023 18:48
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:APIAirflow's REST/HTTP APIarea:CLIarea:Schedulerincluding HA (high availability) schedulerarea:serializationarea:webserverWebserver related Issuesprovider:cncf-kubernetesKubernetes (k8s) provider related issues

Projects

None yet

Development

Successfully merging this pull request may close these issues.

7 participants

@mhenc@ashb@bolkedebruin@potiuk@uranusjr@kosteev@vincbeck