Skip to content

Skip trigger timeout check on occasional db deadlocks - #33172

Closed
vchiapaikeo wants to merge 2 commits into
apache:mainfrom
vchiapaikeo:vchiapaikeo/triggerer-check-deadlocks-v1
Closed

vchiapaikeo wants to merge 2 commits into
apache:mainfrom
vchiapaikeo:vchiapaikeo/triggerer-check-deadlocks-v1

Conversation

@vchiapaikeo

@vchiapaikeo vchiapaikeo commented Aug 7, 2023 •

Copy link
Copy Markdown
Contributor

closes: #32698

We occasionally hit DB deadlocks during the trigger timeout check process. When this happens, the scheduler crashes. See linked issue for graphs, traceback, etc.. I see two potential options:

  1. Retry these transactions
  2. Skip the check entirely and allow another replica or the next timeout check to run

At first, I was leaning toward option 1 but after seeing it in code, I worry that the retries will cause even more deadlocks. That is, if a user set AIRFLOW__DATABASE__MAX_DB_RETRIES to a high number (default is 3) and since AIRFLOW__SCHEDULER__TRIGGER_TIMEOUT_CHECK_INTERVAL already defaults to a very low number, 15s, this could increase deadlock likelihood and cause even more stress on the db.

Therefore, going with 2. Please let me know what you think. This should allow us to avoid crashing the scheduler when deadlocks occur on this (not as critical) transaction.

Manual Testing

Started breeze with the following configurations set:

export AIRFLOW__LOGGING__LOGGING_LEVEL=DEBUG
export AIRFLOW__SCHEDULER__TRIGGER_TIMEOUT_CHECK_INTERVAL=5
image

@vchiapaikeo vchiapaikeo changed the title Retry occasional deadlock on schedulers trigger check process Skip trigger timeout check on occasional db deadlocks Aug 7, 2023
@vchiapaikeo
vchiapaikeo marked this pull request as ready for review August 7, 2023 13:36
@vchiapaikeo
vchiapaikeo force-pushed the vchiapaikeo/triggerer-check-deadlocks-v1 branch from 75ac8b4 to 6317f32 Compare August 7, 2023 13:51
@potiuk

potiuk commented Aug 7, 2023

Copy link
Copy Markdown
Member

I am not sure - but I believe this is not the right fix. With this query the problem is that it is simply not written in the way to grab the right lock on the DagRun while updating task instance table. I think (unlike some other queries and problems we have with deadlocks - this one could be written in the way that it could grab the locks - possibly with SKIP_LOCKED to make it less contentious with scheduler) and avoid the locks rather than reacting to them.

@potiuk

potiuk commented Aug 7, 2023

Copy link
Copy Markdown
Member

Just to add some context: Scheduler only operates on task instances tha belong to DAGRuns that it managed to get "For Update" lock on. This means that any query that modifies a bulk of task instances is bound to hit the deadlock, because it might get a lock on one task instance to update, and then wait for dagrun, while scheduler will do that in reverse direction and will attempt to update the two instance in reverse order.

@vchiapaikeo

Copy link
Copy Markdown
Contributor Author

Hmm yes, that makes sense. Would adding the below condition be sufficient to look at TI's queued by the current replica?

TI.queued_by_job_id == self.job.id,

So then the method would look something like this?

    @provide_session
    def check_trigger_timeouts(self, session: Session = NEW_SESSION) -> None:
        """Mark any "deferred" task as failed if the trigger or execution timeout has passed."""
        self.log.debug("Calling SchedulerJob.check_trigger_timeouts method")

        try:
            num_timed_out_tasks = session.execute(
                update(TI)
                .where(
                    TI.state == TaskInstanceState.DEFERRED,
                    TI.trigger_timeout < timezone.utcnow(),
                    # Only perform update against the ones that were queued by this scheduler
                    TI.queued_by_job_id == self.job.id,
                )
                .values(
                    state=TaskInstanceState.SCHEDULED,
                    next_method="__fail__",
                    next_kwargs={"error": "Trigger/execution timeout"},
                    trigger_id=None,
                )
            ).rowcount

            if num_timed_out_tasks:
                self.log.info("Timed out %i deferred tasks without fired triggers", num_timed_out_tasks)

        except OperationalError as e:
            session.rollback()
            self.log.warning(
                f"Failed to check trigger timeouts due to {e}. Will reattempt at next scheduled check"
            )

The problem I see here is that there may be lingering triggers that do not get cleaned up if replicas get dropped. Maybe this overcomplicates the problem...

@potiuk

potiuk commented Aug 7, 2023 •

Copy link
Copy Markdown
Member

Hmm yes, that makes sense. Would adding the below condition be sufficient to look at TI's queued by the current replica?

TI.queued_by_job_id == self.job.id,

So then the method would look something like this?

    @provide_session
    def check_trigger_timeouts(self, session: Session = NEW_SESSION) -> None:
        """Mark any "deferred" task as failed if the trigger or execution timeout has passed."""
        self.log.debug("Calling SchedulerJob.check_trigger_timeouts method")

        try:
            num_timed_out_tasks = session.execute(
                update(TI)
                .where(
                    TI.state == TaskInstanceState.DEFERRED,
                    TI.trigger_timeout < timezone.utcnow(),
                    # Only perform update against the ones that were queued by this scheduler
                    TI.queued_by_job_id == self.job.id,
                )
                .values(
                    state=TaskInstanceState.SCHEDULED,
                    next_method="__fail__",
                    next_kwargs={"error": "Trigger/execution timeout"},
                    trigger_id=None,
                )
            ).rowcount

            if num_timed_out_tasks:
                self.log.info("Timed out %i deferred tasks without fired triggers", num_timed_out_tasks)

        except OperationalError as e:
            session.rollback()
            self.log.warning(
                f"Failed to check trigger timeouts due to {e}. Will reattempt at next scheduled check"
            )

The problem I see here is that there may be lingering triggers that do not get cleaned up if replicas get dropped. Maybe this overcomplicates the problem...

Nope. It's different. IMHO you should not limit it by job_id, but you should add dag_run "FOR UPDATE" section with SKIP_LOCKED condition - basically to join the task_instance with dag_run they belong to, and make sure that dag_run is locked "for update". I am not super expert in the sqlalchemy queries to be able to tell exactly how it should be done, my knowledge is more based on the "relational DB knowledge" I have, and learning from @ashb's https://www.youtube.com/watch?v=DYC4-xElccE talk on how scheduler works internally - but there are quite a few such sqlalchemy queries in Airflow code that already do it I believe.

Maybe others who know better could chime in as well?

@vijayasarathib

This comment was marked as off-topic.

@vchiapaikeo

Copy link
Copy Markdown
Contributor Author

So it turns out that this issue actually stemmed from something completely unrelated. Our team will post a discussion / issue about that later. Long story short, we were using AIRFLOW__DATABASE__SQL_ALCHEMY_CONN_CMD instead of AIRFLOW__DATABASE__SQL_ALCHEMY_CONN to define our database connection string. A trace showed that while making a bulk fetch query (a simple .all() on the trigger table), we are continually making calls to configuration.py to subprocess out for the connection string for EVERY record that is returned checking whether or not the database supported json 😭 . Switching to defining the connection with AIRFLOW__DATABASE__SQL_ALCHEMY_CONN resolved our issues and gave us huge performance boosts on both our scheduler and triggerer processes.

   65686 function calls (64626 primitive calls) in 29.928 seconds                          
                                                                                
   Ordered by: cumulative time                                                                                                      
                                                                                                                                        
   ncalls  tottime  percall  cumtime  percall filename:lineno(function)                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                           [506/45463]
        1    0.000    0.000   29.928   29.928 /home/airflow/.local/lib/python3.10/site-packages/sqlalchemy/orm/query.py:2757(all)               
        1    0.000    0.000   29.923   29.923 /home/airflow/.local/lib/python3.10/site-packages/sqlalchemy/engine/result.py:1468(all)      
        1    0.000    0.000   29.923   29.923 /home/airflow/.local/lib/python3.10/site-packages/sqlalchemy/engine/result.py:395(_allrows)  
        1    0.000    0.000   29.923   29.923 /home/airflow/.local/lib/python3.10/site-packages/sqlalchemy/engine/result.py:1388(_fetchall_impl)       
        1    0.000    0.000   29.923   29.923 /home/airflow/.local/lib/python3.10/site-packages/sqlalchemy/engine/result.py:1808(_fetchall_impl)                                                                                                                                                                                                                                                                                                                                                                                
        2    0.000    0.000   29.922   14.961 /home/airflow/.local/lib/python3.10/site-packages/sqlalchemy/orm/loading.py:135(chunks)
        1    0.000    0.000   29.921   29.921 /home/airflow/.local/lib/python3.10/site-packages/sqlalchemy/engine/result.py:390(_raw_all_rows)
        1    0.001    0.001   29.921   29.921 /home/airflow/.local/lib/python3.10/site-packages/sqlalchemy/engine/result.py:393(<listcomp>)                                                                                                                                                                                                                                                                                                                                                                                                                        
      125    0.000    0.000   29.919    0.239 /home/airflow/.local/lib/python3.10/site-packages/sqlalchemy/sql/type_api.py:1711(process) 
      125    0.002    0.000   29.915    0.239 /home/airflow/.local/lib/python3.10/site-packages/airflow/utils/sqlalchemy.py:146(process_result_value)  
      125    0.001    0.000   29.909    0.239 /home/airflow/.local/lib/python3.10/site-packages/airflow/utils/sqlalchemy.py:122(db_supports_json)    
      125    0.001    0.000   29.908    0.239 /home/airflow/.local/lib/python3.10/site-packages/airflow/configuration.py:562(get)                  
      125    0.000    0.000   29.907    0.239 /home/airflow/.local/lib/python3.10/site-packages/airflow/configuration.py:732(_get_environment_variables)
      125    0.002    0.000   29.907    0.239 /home/airflow/.local/lib/python3.10/site-packages/airflow/configuration.py:478(_get_env_var_option)
      125    0.002    0.000   29.902    0.239 /home/airflow/.local/lib/python3.10/site-packages/airflow/configuration.py:103(run_command)                
      125    0.001    0.000   29.786    0.238 /usr/local/lib/python3.10/subprocess.py:1110(communicate)                                                       
      125    0.006    0.000   29.785    0.238 /usr/local/lib/python3.10/subprocess.py:1952(_communicate)                                             
      250    0.003    0.000   29.762    0.119 /usr/local/lib/python3.10/selectors.py:403(select)                                         
      250   29.758    0.119   29.758    0.119 {method 'poll' of 'select.poll' objects}                                                                  
      125    0.002    0.000    0.100    0.001 /usr/local/lib/python3.10/subprocess.py:758(__init__)                                                                  
      125    0.004    0.000    0.094    0.001 /usr/local/lib/python3.10/subprocess.py:1687(_execute_child)                                             
      125    0.069    0.001    0.069    0.001 {built-in method _posixsubprocess.fork_exec}                                            
      125    0.001    0.000    0.013    0.000 /usr/local/lib/python3.10/shlex.py:305(split)                                                 
      500    0.001    0.000    0.010    0.000 /usr/local/lib/python3.10/shlex.py:299(__next__)                                  
      500    0.001    0.000    0.010    0.000 /usr/local/lib/python3.10/shlex.py:101(get_token)                                            
      500    0.007    0.000    0.009    0.000 /usr/local/lib/python3.10/shlex.py:133(read_token)                                                                                                                                                                                                                                                                                                                                                                                                                                
     1375    0.001    0.000    0.006    0.000 /usr/local/lib/python3.10/subprocess.py:1775(<genexpr>)                                  
      125    0.002    0.000    0.006    0.000 /usr/local/lib/python3.10/os.py:620(get_exec_path)                                                            
      250    0.000    0.000    0.005    0.000 /usr/local/lib/python3.10/subprocess.py:1204(wait)                                                                                                                                                                                                                                                                                                                                                                                                                                                                   
        1    0.000    0.000    0.005    0.005 /home/airflow/.local/lib/python3.10/site-packages/sqlalchemy/orm/query.py:2911(_iter) 
        1    0.000    0.000    0.005    0.005 /home/airflow/.local/lib/python3.10/site-packages/sqlalchemy/orm/session.py:1563(execute)                  
      250    0.001    0.000    0.005    0.000 /usr/local/lib/python3.10/subprocess.py:1911(_wait)                                                             
      125    0.001    0.000    0.004    0.000 /usr/local/lib/python3.10/subprocess.py:1227(_close_pipe_fds)                                          
 1057/125    0.002    0.000    0.004    0.000 /home/airflow/.local/lib/python3.10/site-packages/airflow/serialization/serialized_objects.py:458(deserialize)
        1    0.000    0.000    0.004    0.004 /home/airflow/.local/lib/python3.10/site-packages/sqlalchemy/engine/base.py:1691(_execute_20)          
        1    0.000    0.000    0.004    0.004 /home/airflow/.local/lib/python3.10/site-packages/sqlalchemy/sql/elements.py:330(_execute_on_connection)   
        1    0.000    0.000    0.004    0.004 /home/airflow/.local/lib/python3.10/site-packages/sqlalchemy/engine/base.py:1523(_execute_clauseelement)        
      125    0.000    0.000    0.004    0.000 /usr/local/lib/python3.10/subprocess.py:1898(_try_wait)                                                
     1250    0.003    0.000    0.004    0.000 /usr/local/lib/python3.10/posixpath.py:71(join)                                            
        1    0.000    0.000    0.004    0.004 /home/airflow/.local/lib/python3.10/site-packages/sqlalchemy/engine/base.py:1768(_execute_context)           
      125    0.000    0.000    0.004    0.000 /home/airflow/.local/lib/python3.10/site-packages/sqlalchemy/sql/sqltypes.py:2675(process)    
      125    0.004    0.000    0.004    0.000 {built-in method posix.waitpid}                                                                                                                                                                                                                                                                                                                                                                           
      125    0.001    0.000    0.003    0.000 /usr/local/lib/python3.10/json/__init__.py:299(loads)                                                                                                                                                                                                                                                                                                                            
        1    0.000    0.000    0.003    0.003 /home/airflow/.local/lib/python3.10/site-packages/sqlalchemy/engine/default.py:735(do_execute)                                                                                                                                                                                                                                                                                   
        1    0.000    0.000    0.003    0.003 /home/airflow/.local/lib/python3.10/site-packages/MySQLdb/cursors.py:171(execute)                                                                                                                                                                                                                                                                                                       
      250    0.001    0.000    0.003    0.000 /usr/local/lib/python3.10/selectors.py:366(unregister)                                                                                                                                                                                                                                                                                                                                  
      625    0.001    0.000    0.003    0.000 /usr/local/lib/python3.10/os.py:675(__getitem__)                                                                                                                                                                                                                                                                                                                                                     
      250    0.001    0.000    0.003    0.000 /usr/local/lib/python3.10/selectors.py:352(register)                                                                                                                                                                                                                                                                                                                                                  
      125    0.001    0.000    0.003    0.000 /usr/local/lib/python3.10/json/decoder.py:332(decode)                                                                                                                                                                                                                                                                                                                                   
  240/125    0.001    0.000    0.003    0.000 /home/airflow/.local/lib/python3.10/site-packages/airflow/serialization/serialized_objects.py:476(<dictcomp>)                                                                                                                                                                                                                                                                                        
[...truncated]

Relevant callsites:

Triggerer Model:
https://github.com/apache/airflow/blob/2.5.3/airflow/models/trigger.py#L57

@potiuk

potiuk commented Aug 17, 2023

Copy link
Copy Markdown
Member

That's a very interesting finding. Looks like it needs fixing indeed.

@potiuk

potiuk commented Aug 17, 2023

Copy link
Copy Markdown
Member

Ping me on an issue when you create it @vchiapaikeo.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Triggerer Timeout Checks Cause Deadlocks when Multiple Scheduler Replicas Running at Peak

3 participants