DUE TO SPAM, SIGN-UP IS DISABLED. Goto Selfserve wiki signup and request an account.
Status
| Page properties | |||||||||||||||||||
|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|
|
Motivation
Currently Airflow requires DAG files to be present on a file system that is accessible to the scheduler, webserver, and workers. Given that more and more people are running airflow in a distributed setup to achieve higher scalability, it becomes more and more difficult to guarantee a file system that is accessible and synchronized amongst services. By allowing Airflow to fetch DAG files from a remote source outside the file system local to the service, this grant a much greater flexibility, eases implementation, and standardizes ways to sync remote sources of DAGs with Airflow.
Proposed Solutions
Option 1: DAG Repository (short term)
DAGs are persisted in remote filesystem-like storage and Airflow need
Suggested Implementation
Since DAG files are stored remotely and Airflow needs to know where to find them, . DAG manifestrepository is introduced to record the location remote root directory of the DAG files. Airflow looks at the DAG manifest and fetches DAG files from remote storage using DagFetcher to local file system. Airflow caches the DAGs on the local filesystem and will only fetch from remote if the local copy is stale.
Dag Manifest:
Currently Airflow assumes valid DAGs are any Dag object that it finds in the python files in thePrior to DAG loading, Airflow would download files from the remote DAG repository and cache it under the local filesystem directory under $AIRFLOW_HOME/dags
directory. However, to support Airflow fetching DAG files from remote sources, we need a figure out a way to record the remote location of each DAG, which is the DAG manifest. The implementation of DAG manifest can simply be a manifest.yml on the filesystem or we can store the DAG manifest on a database..
DAG Repository
We will create a file, remote_repositories.json, to record the root directory of DAGs on the remote storage system. Multiple repositories are supported.
Format
| Code Block |
|---|
dag_repositories: [
"repo1": {
"url": " |
Format
The Dag Manifest would be composed of manifest entries, an entry would be defined as follows:
| Code Block | ||
|---|---|---|
| ||
dag_manifest_entry:
dag_id:
uri: where dag can be found
conn_id: connection id to use to interact with remote location |
File-based DAG manifest
Airflow services will look at $AIRFLOW_HOME/manifest.yml for the DAG manifest. The manifest.yml contains all the DAG entries. We should expect a manifest.yml like:
| Code Block |
|---|
dag_id_1 - uri: s3://my-bucket/dag_id_1_1.zip - dags", "conn_id": some conn id dag_id_2 - uri: "blahblah" }, "repo2": { "url": "git://repo_name/dags", "conn_id": "blahblah2" } ] |
DagFetcher
The following is the DagFetcher interface, we will implement different fetchers for different storage system, GitDagFetcher and S3DagFetcher. Say we have a remote_repositories.json configuration like above. DagFetcher would download files under s3://my-bucket/dags to $AIRFLOW_HOME/dags/repo_id/
| Code Block |
|---|
class BaseDagFetcher(): def fetch(repo_id, url,/dag_id_2_3.zip - conn_id, file_path=None): some conn id |
Custom DAG manifest
The manifest can also be generated by a callable supplied in the airflow.cfg that when called would generate a list of entries i.e
| Code Block | ||
|---|---|---|
| ||
[core]
# callable to fetch dag manifest list
dag_manifest_entries = my_config.get_dag_manifest_entries |
The DAG manifest can be stored on S3 and my_config.get_dag_manifest_entries will read the manifest from S3.
Backwards Compatibility
To maintain backward compatibility, we can support a migration script that crawl the DAGs in $AIRFLOE_HOME and populates a manifest.yml file.
Benefits
- With the manifest people are able to more explicitly note which dags should be looked at for by Airflow
- Airflow no longer has to crawl through a directory importing various files possibly causing problems
- Users are not forced to allow for a way to crawl various remote sources
- Allowing listing the connection id makes it easy to have multiple remote dag locations
- We can get rid of
, which requires strings such as "airflow" and "DAG" to be present in DAG file.Jira server ASF JIRA serverId 5aa69414-a9e9-3523-82ec-879b028fb15b key AIRFLOW-97
DAG URIs:
DAG locations will be given via URI, i.e.s3://my-bucket/dag1.zip, local:////dags/day1.zip
"""
Download files from remote storage to local directory under $AIRFLOW_HOME/dags/repo_id
""" |
Proposed changes
DagBag
We should ensure that we are loading the latest DAGs cache copy, thus we should fetch DAGs from remote repo before we load the DagBag.
Scheduler
Currently, DagFileProcessorManager periodically calls DagFileProcessorManager._refresh_dag_dir to look for new DAG files. We should change this method to fetch DAGs from remote at first .
Versioning:
DAG URI is also DAG version.
- We load the DAG from the same URI throughout the entire DAG run even if the DAG manifest was changed to a new DAG URI. We will add a URI attribute to the DagRun model to persist the URI used for each DAG run.
- Users are free to define their own URI naming convention.
- Same version/URI should not be re-used.
- We load the DAG from the same URI throughout the entire DAG run even if the DAG manifest was changed to a new DAG URI. We will add a URI attribute to the DagRun model to persist the URI used for each DAG run.
Caching:
Moving DAGs to a remote location will introduce network overhead so we should cache DAGs and avoid unnecessary fetch. We should only re-fetch a DAG when the URI in the DAG manifest is changed (or we find that the last_update_time of the DAG file is changed).
Cache Location
In order to avoid making a remote fetch every time the dag needs to be run it will be best to keep a local cache of dag files for individual Airflow services to use
- User will specify a cached location. Each service (scheduler, webserver, and worker) will cache the latest DAG file files in the cache location.
- We will store DAGs in the cached location in the following structure: cached_dags_path/<DAG_ID>/<DAG_URI>
- Users should not modify anything under cached_dags_path/
Dag Fetching
- We also save the last_modified_date of the remote files.
Proposed changes
DagBag
DagBag ChangesThe main part of the code base that we will need to change is in DagBag. In collect_dags, we will go through each entries in defined in the DAG manifest and download the DAG files if cache is invalid. In download_dag_file_and_add_to_cache, we will use different fetching implementation for different kind of uri, e.g. s3 or git.
| Code Block | ||
|---|---|---|
| ||
class DagBag(): def collect_dags(): for entry in get_dag_manifest_entries(): if entry.uri is stored locally: self.process_file(entry.uri, only_if_updated=True) continue # the DAG is stored remotely dag_cache_path = self.get_cache_path(entry.dag_id, entry.uri) if os.direxists(dag_cache_path) and self.cached_dag_file_is_latest(dag_id, uri, dag_cache_path): # we have the latest cache of DAG continue else: download_dag_file_and_add_to_cache(dag_id, entry.uri) self.process_file(cache_path) def get_last_modified_date_from_local(dag_id, uri, dag_cache_path): pass def cached_dag_file_is_latest(dag_id, uri, dag_cache_path): uri_type = get_uri_type(uri) dag_fetcher = self.get_dag_fetcher(uri_type) return self.get_last_modified_date_from_local(dag_id, uri, dag_cache_path) == dag_fetcher.get_last_modified_date(dag_id, uri)""" Check if the DAG file on remote storage is changed since the last download. """ def download_dag_file_and_add_to_cache(dag_id, uri, dag_cache_path): """ Download DAG files from remote location to the local dag_cache_path. """ uri_type = get_uri_type(uri) dag_fetcher = self.get_dag_fetcher(uri_type) dag_fetcher.fetch_and_cache_dag(dag_id, uri, dag_cache_path) def get_cache_path(dag_id, uri): return cached_dags_path + "/" + dag_id + "/" + uri class DagFetcher(): def fetch_and_cache_dag(dag_id, uri, conn_id, dag_cache_path): """ Download the DAG file from remote to local file system under dag_cache_path and record the last_modified_date in a file on local filesystem. """ raise NotImplementedError() def get_last_modified_date(dag_id, uri, conn_id) """ When the DAG is last modified on remote system. """ raise NotImplementedError() |
Scheduler
ChangesCurrently the scheduler checks what dags are on disk by calling list_py_file_paths this will need to be changed to instead look at the manifest as we can no longer crawl the file system, and instead crawl manifest entries
and load the files on the local filesystem cache. Airflow scheduler will need to persist the URI into the DagRun table when it creates a new DagRun.
DagRun versioning
If we implement versioning of dags it will require a number of changes to the current scheduler. The biggest issue comes from how the scheduler currently propagates the dag object to it's various function calls for task scheduling. As is the scheduler loads in the dag objects that are found in the filesystem, and these passed along to the resulting functions. In order to implement versions we would need to associate a certain dag version/uri to a DagRun, when previous DagRuns are fetched we'll need to check if they were for an earlier dag version/uri and fetch that version/uri if necessary.
Here is the exact part of _process_task_instances where this check would need to happen
| language | py |
|---|
.
""" # update the state of the previously active dag runs dag_runs = DagRun.find(dag_id=dag.dag_id, state=State.RUNNING, session=session) active_dag_runs = [] for run in dag_runs: self.log.info("Examining DAG run %s", run) # don't consider runs that are executed in the future if run.execution_date > timezone.utcnow(): self.log.error( "Execution date is in future: %s", run.execution_date ) continue if len(active_dag_runs) >= dag.max_active_runs: self.log.info("Active dag runs > max_active_run.") continue # skip backfill dagruns for now as long as they are not really scheduled if run.is_backfill: continue # todo: run.dag is transient but needs to be set run.dag = dagIn this case the dag object will be the object loaded in from the current version listed in the manifest, this last line here should check to pin run.dag to be equal to the dag version of the DagRun
DAG Model
The main thing we'll have to make sure stays consistent about the dag model is the thefileloc attribute points to a file location on disk accessible to the webserver, one way to do this is to set fileloc to be the local cache location, and that when getting a dag from the DagBag we ensure the cache file is available on disk