diff --git a/.flake8 b/.flake8 index 14de564a3291e..01f648b2143c6 100644 --- a/.flake8 +++ b/.flake8 @@ -6,3 +6,4 @@ format = ${cyan}%(path)s${reset}:${yellow_bold}%(row)d${reset}:${green_bold}%(co per-file-ignores = airflow/models/__init__.py:F401 airflow/models/sqla_models.py:F401 + airflow/config_templates/built_in_defaults:E501 diff --git a/.pre-commit-config.yaml b/.pre-commit-config.yaml index 1a67aaf66ad26..a033f921a0003 100644 --- a/.pre-commit-config.yaml +++ b/.pre-commit-config.yaml @@ -583,7 +583,7 @@ repos: files: config\.yml$|default_airflow\.cfg$|default\.cfg$ pass_filenames: false require_serial: true - additional_dependencies: ['pyyaml'] + additional_dependencies: ['pyyaml', 'jinja2', 'black==22.3.0', 'rich'] - id: check-boring-cyborg-configuration name: Checks for Boring Cyborg configuration consistency language: python @@ -894,14 +894,14 @@ repos: files: \.py$ exclude: ^provider_packages|^docs|^airflow/_vendor/|^airflow/providers|^airflow/migrations|^dev|^tests/system/providers|^tests/providers require_serial: true - additional_dependencies: ['rich>=12.4.4', 'inputimeout'] + additional_dependencies: ['rich>=12.4.4', 'inputimeout', 'jinja2'] - id: run-mypy name: Run mypy for providers language: python entry: ./scripts/ci/pre_commit/pre_commit_mypy.py --namespace-packages files: ^airflow/providers/.*\.py$|^tests/system/providers/\*.py|^tests/providers/\*.py require_serial: true - additional_dependencies: ['rich>=12.4.4', 'inputimeout'] + additional_dependencies: ['rich>=12.4.4', 'inputimeout', 'jinja2'] - id: run-mypy name: Run mypy for /docs/ folder language: python @@ -909,7 +909,7 @@ repos: files: ^docs/.*\.py$ exclude: ^docs/rtd-deprecation require_serial: true - additional_dependencies: ['rich>=12.4.4', 'inputimeout'] + additional_dependencies: ['rich>=12.4.4', 'inputimeout', 'jinja2'] - id: run-flake8 name: Run flake8 language: python @@ -917,7 +917,7 @@ repos: files: \.py$ pass_filenames: true exclude: ^airflow/_vendor/ - additional_dependencies: ['rich>=12.4.4', 'inputimeout'] + additional_dependencies: ['rich>=12.4.4', 'inputimeout', 'jinja2'] - id: update-migration-references name: Update migration ref doc language: python diff --git a/airflow/config_templates/built_in_defaults.py b/airflow/config_templates/built_in_defaults.py new file mode 100644 index 0000000000000..27779a8ab7984 --- /dev/null +++ b/airflow/config_templates/built_in_defaults.py @@ -0,0 +1,338 @@ +# +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you under the Apache License, Version 2.0 (the +# "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, +# software distributed under the License is distributed on an +# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +# KIND, either express or implied. See the License for the +# specific language governing permissions and limitations +# under the License. + +# Those are built-in defaults for Airflow's default configuration. When there +# are no defaults stored in airflow configuration file, those defaults are used +# as a fallback for the configuration values +# +# PLEASE DO NOT MODIFY THIS FILE IT IS AUTOGENERATED FROM THE CONFIGURATION.YAML FILE +from __future__ import annotations + +built_in_defaults: dict[str, dict[str, str]] = { + "core": { + "dags_folder": "{AIRFLOW_HOME}/dags", + "hostname_callable": "airflow.utils.net.getfqdn", + "default_timezone": "utc", + "executor": "SequentialExecutor", + "parallelism": "32", + "max_active_tasks_per_dag": "16", + "dags_are_paused_at_creation": "True", + "max_active_runs_per_dag": "16", + "load_examples": "True", + "plugins_folder": "{AIRFLOW_HOME}/plugins", + "execute_tasks_new_python_interpreter": "False", + "fernet_key": "{FERNET_KEY}", + "donot_pickle": "True", + "dagbag_import_timeout": "30.0", + "dagbag_import_error_tracebacks": "True", + "dagbag_import_error_traceback_depth": "2", + "dag_file_processor_timeout": "50", + "task_runner": "StandardTaskRunner", + "default_impersonation": "", + "security": "", + "unit_test_mode": "False", + "enable_xcom_pickling": "False", + "allowed_deserialization_classes": "airflow\\..*", + "killed_task_cleanup_time": "60", + "dag_run_conf_overrides_params": "True", + "dag_discovery_safe_mode": "True", + "dag_ignore_file_syntax": "regexp", + "default_task_retries": "0", + "default_task_retry_delay": "300", + "max_task_retry_delay": "86400", + "default_task_weight_rule": "downstream", + "default_task_execution_timeout": "", + "min_serialized_dag_update_interval": "30", + "compress_serialized_dags": "False", + "min_serialized_dag_fetch_interval": "10", + "max_num_rendered_ti_fields_per_task": "30", + "check_slas": "True", + "xcom_backend": "airflow.models.xcom.BaseXCom", + "lazy_load_plugins": "True", + "lazy_discover_providers": "True", + "hide_sensitive_var_conn_fields": "True", + "sensitive_var_conn_names": "", + "default_pool_task_slot_count": "128", + "max_map_length": "1024", + "daemon_umask": "0o077", + "database_access_isolation": "False", + }, + "database": { + "sql_alchemy_conn": "sqlite:///{AIRFLOW_HOME}/airflow.db", + "sql_engine_encoding": "utf-8", + "sql_alchemy_pool_enabled": "True", + "sql_alchemy_pool_size": "5", + "sql_alchemy_max_overflow": "10", + "sql_alchemy_pool_recycle": "1800", + "sql_alchemy_pool_pre_ping": "True", + "sql_alchemy_schema": "", + "load_default_connections": "True", + "max_db_retries": "3", + }, + "logging": { + "base_log_folder": "{AIRFLOW_HOME}/logs", + "remote_logging": "False", + "remote_log_conn_id": "", + "google_key_path": "", + "remote_base_log_folder": "", + "encrypt_s3_logs": "False", + "logging_level": "INFO", + "celery_logging_level": "", + "fab_logging_level": "WARNING", + "logging_config_class": "", + "colored_console_log": "True", + "colored_log_format": ( + "[%%(blue)s%%(asctime)s%%(reset)s] {{%%(blue)s%%(filename)s:%%(reset)s%%(lineno)d}}" + " %%(log_color)s%%(levelname)s%%(reset)s - %%(log_color)s%%(message)s%%(reset)s" + ), + "colored_formatter_class": "airflow.utils.log.colored_log.CustomTTYColoredFormatter", + "log_format": "[%%(asctime)s] {{%%(filename)s:%%(lineno)d}} %%(levelname)s - %%(message)s", + "simple_log_format": "%%(asctime)s %%(levelname)s - %%(message)s", + "dag_processor_log_target": "file", + "dag_processor_log_format": ( + "[%%(asctime)s] [SOURCE:DAG_PROCESSOR] {{%%(filename)s:%%(lineno)d}} %%(levelname)s -" + " %%(message)s" + ), + "log_formatter_class": "airflow.utils.log.timezone_aware.TimezoneAware", + "secret_mask_adapter": "", + "task_log_prefix_template": "", + "log_filename_template": ( + "dag_id={{ ti.dag_id }}/run_id={{ ti.run_id }}/task_id={{ ti.task_id }}/{%% if ti.map_index >= 0" + " %%}map_index={{ ti.map_index }}/{%% endif %%}attempt={{ try_number }}.log" + ), + "log_processor_filename_template": "{{ filename }}.log", + "dag_processor_manager_log_location": ( + "{AIRFLOW_HOME}/logs/dag_processor_manager/dag_processor_manager.log" + ), + "task_log_reader": "task", + "extra_logger_names": "", + "worker_log_server_port": "8793", + }, + "metrics": { + "statsd_on": "False", + "statsd_host": "localhost", + "statsd_port": "8125", + "statsd_prefix": "airflow", + "statsd_allow_list": "", + "stat_name_handler": "", + "statsd_datadog_enabled": "False", + "statsd_datadog_tags": "", + }, + "secrets": {"backend": "", "backend_kwargs": ""}, + "cli": {"api_client": "airflow.api.client.local_client", "endpoint_url": "http://localhost:8080"}, + "debug": {"fail_fast": "False"}, + "api": { + "enable_experimental_api": "False", + "auth_backends": "airflow.api.auth.backend.session", + "maximum_page_limit": "100", + "fallback_page_limit": "100", + "google_oauth2_audience": "", + "google_key_path": "", + "access_control_allow_headers": "", + "access_control_allow_methods": "", + "access_control_allow_origins": "", + }, + "lineage": {"backend": ""}, + "atlas": {"sasl_enabled": "False", "host": "", "port": "21000", "username": "", "password": ""}, + "operators": { + "default_owner": "airflow", + "default_cpus": "1", + "default_ram": "512", + "default_disk": "512", + "default_gpus": "0", + "default_queue": "default", + "allow_illegal_arguments": "False", + }, + "hive": {"default_hive_mapred_queue": ""}, + "webserver": { + "base_url": "http://localhost:8080", + "default_ui_timezone": "UTC", + "web_server_host": "0.0.0.0", + "web_server_port": "8080", + "web_server_ssl_cert": "", + "web_server_ssl_key": "", + "session_backend": "database", + "web_server_master_timeout": "120", + "web_server_worker_timeout": "120", + "worker_refresh_batch_size": "1", + "worker_refresh_interval": "6000", + "reload_on_plugin_change": "False", + "secret_key": "{SECRET_KEY}", + "workers": "4", + "worker_class": "sync", + "access_logfile": "-", + "error_logfile": "-", + "access_logformat": "", + "expose_config": "False", + "expose_hostname": "True", + "expose_stacktrace": "False", + "dag_default_view": "grid", + "dag_orientation": "LR", + "log_fetch_timeout_sec": "5", + "log_fetch_delay_sec": "2", + "log_auto_tailing_offset": "30", + "log_animation_speed": "1000", + "hide_paused_dags_by_default": "False", + "page_size": "100", + "navbar_color": "#fff", + "default_dag_run_display_number": "25", + "enable_proxy_fix": "False", + "proxy_fix_x_for": "1", + "proxy_fix_x_proto": "1", + "proxy_fix_x_host": "1", + "proxy_fix_x_port": "1", + "proxy_fix_x_prefix": "1", + "cookie_secure": "False", + "cookie_samesite": "Lax", + "default_wrap": "False", + "x_frame_enabled": "True", + "show_recent_stats_for_completed_runs": "True", + "update_fab_perms": "True", + "session_lifetime_minutes": "43200", + "instance_name_has_markup": "False", + "auto_refresh_interval": "3", + "warn_deployment_exposure": "True", + "audit_view_excluded_events": "gantt,landing_times,tries,duration,calendar,graph,grid,tree,tree_data", + "enable_swagger_ui": "True", + "run_internal_api": "False", + }, + "email": { + "email_backend": "airflow.utils.email.send_email_smtp", + "email_conn_id": "smtp_default", + "default_email_on_retry": "True", + "default_email_on_failure": "True", + }, + "smtp": { + "smtp_host": "localhost", + "smtp_starttls": "True", + "smtp_ssl": "False", + "smtp_port": "25", + "smtp_mail_from": "airflow@example.com", + "smtp_timeout": "30", + "smtp_retry_limit": "5", + }, + "sentry": {"sentry_on": "false", "sentry_dsn": ""}, + "local_kubernetes_executor": {"kubernetes_queue": "kubernetes"}, + "celery_kubernetes_executor": {"kubernetes_queue": "kubernetes"}, + "celery": { + "celery_app_name": "airflow.executors.celery_executor", + "worker_concurrency": "16", + "worker_prefetch_multiplier": "1", + "worker_enable_remote_control": "true", + "broker_url": "redis://redis:6379/0", + "flower_host": "0.0.0.0", + "flower_url_prefix": "", + "flower_port": "5555", + "flower_basic_auth": "", + "sync_parallelism": "0", + "celery_config_options": "airflow.config_templates.default_celery.DEFAULT_CELERY_CONFIG", + "ssl_active": "False", + "ssl_key": "", + "ssl_cert": "", + "ssl_cacert": "", + "pool": "prefork", + "operation_timeout": "1.0", + "task_track_started": "True", + "task_adoption_timeout": "600", + "stalled_task_timeout": "0", + "task_publish_max_retries": "3", + "worker_precheck": "False", + }, + "celery_broker_transport_options": {}, + "dask": {"cluster_address": "127.0.0.1:8786", "tls_ca": "", "tls_cert": "", "tls_key": ""}, + "scheduler": { + "job_heartbeat_sec": "5", + "scheduler_heartbeat_sec": "5", + "num_runs": "-1", + "scheduler_idle_sleep_time": "1", + "min_file_process_interval": "30", + "parsing_cleanup_interval": "60", + "dag_dir_list_interval": "300", + "print_stats_interval": "30", + "pool_metrics_interval": "5.0", + "scheduler_health_check_threshold": "30", + "enable_health_check": "False", + "scheduler_health_check_server_port": "8974", + "orphaned_tasks_check_interval": "300.0", + "child_process_log_directory": "{AIRFLOW_HOME}/logs/scheduler", + "scheduler_zombie_task_threshold": "300", + "zombie_detection_interval": "10.0", + "catchup_by_default": "True", + "ignore_first_depends_on_past_by_default": "True", + "max_tis_per_query": "512", + "use_row_level_locking": "True", + "max_dagruns_to_create_per_loop": "10", + "max_dagruns_per_loop_to_schedule": "20", + "schedule_after_task_execution": "True", + "parsing_processes": "2", + "file_parsing_sort_mode": "modified_time", + "standalone_dag_processor": "False", + "max_callbacks_per_loop": "20", + "dag_stale_not_seen_duration": "600", + "use_job_schedule": "True", + "allow_trigger_in_future": "False", + "trigger_timeout_check_interval": "15", + }, + "triggerer": {"default_capacity": "1000"}, + "kerberos": { + "ccache": "/tmp/airflow_krb5_ccache", + "principal": "airflow", + "reinit_frequency": "3600", + "kinit_path": "kinit", + "keytab": "airflow.keytab", + "forwardable": "True", + "include_ip": "True", + }, + "elasticsearch": { + "host": "", + "log_id_template": "{{dag_id}}-{{task_id}}-{{run_id}}-{{map_index}}-{{try_number}}", + "end_of_log_mark": "end_of_log", + "frontend": "", + "write_stdout": "False", + "json_format": "False", + "json_fields": "asctime, filename, lineno, levelname, message", + "host_field": "host", + "offset_field": "offset", + "index_patterns": "_all", + }, + "elasticsearch_configs": {"use_ssl": "False", "verify_certs": "True"}, + "kubernetes_executor": { + "pod_template_file": "", + "worker_container_repository": "", + "worker_container_tag": "", + "namespace": "default", + "delete_worker_pods": "True", + "delete_worker_pods_on_failure": "False", + "worker_pods_creation_batch_size": "1", + "multi_namespace_mode": "False", + "multi_namespace_mode_namespace_list": "", + "in_cluster": "True", + "kube_client_request_args": "", + "delete_option_kwargs": "", + "enable_tcp_keepalive": "True", + "tcp_keep_idle": "120", + "tcp_keep_intvl": "30", + "tcp_keep_cnt": "6", + "verify_ssl": "True", + "worker_pods_pending_timeout": "300", + "worker_pods_pending_timeout_check_interval": "120", + "worker_pods_queued_check_interval": "60", + "worker_pods_pending_timeout_batch_size": "100", + }, + "sensors": {"default_timeout": "604800"}, +} diff --git a/airflow/config_templates/default_airflow.cfg b/airflow/config_templates/default_airflow.cfg index 4bd2883563f11..caa89c90cd7d1 100644 --- a/airflow/config_templates/default_airflow.cfg +++ b/airflow/config_templates/default_airflow.cfg @@ -30,7 +30,7 @@ [core] # The folder where your airflow pipelines live, most likely a # subfolder in a code repository. This path must be absolute. -dags_folder = {AIRFLOW_HOME}/dags +# dags_folder = {AIRFLOW_HOME}/dags # Hostname by providing a path to a callable, which will resolve the hostname. # The format is "package.function". @@ -40,23 +40,23 @@ dags_folder = {AIRFLOW_HOME}/dags # # No argument should be required in the function specified. # If using IP address as hostname is preferred, use value ``airflow.utils.net.get_host_ip_address`` -hostname_callable = airflow.utils.net.getfqdn +# hostname_callable = airflow.utils.net.getfqdn # Default timezone in case supplied date times are naive # can be utc (default), system, or any IANA timezone string (e.g. Europe/Amsterdam) -default_timezone = utc +# default_timezone = utc # The executor class that airflow should use. Choices include # ``SequentialExecutor``, ``LocalExecutor``, ``CeleryExecutor``, ``DaskExecutor``, # ``KubernetesExecutor``, ``CeleryKubernetesExecutor`` or the # full import path to the class when using a custom executor. -executor = SequentialExecutor +# executor = SequentialExecutor # This defines the maximum number of task instances that can run concurrently per scheduler in # Airflow, regardless of the worker count. Generally this value, multiplied by the number of # schedulers in your cluster, is the maximum number of task instances with the running # state in the metadata database. -parallelism = 32 +# parallelism = 32 # The maximum number of task instances allowed to run concurrently in each DAG. To calculate # the number of tasks that is running concurrently for a DAG, add up the number of running @@ -65,15 +65,15 @@ parallelism = 32 # # An example scenario when this would be useful is when you want to stop a new dag with an early # start date from stealing all the executor slots in a cluster. -max_active_tasks_per_dag = 16 +# max_active_tasks_per_dag = 16 # Are DAGs paused by default at creation -dags_are_paused_at_creation = True +# dags_are_paused_at_creation = True # The maximum number of active DAG runs per DAG. The scheduler will not create more DAG runs # if it reaches the limit. This is configurable at the DAG level with ``max_active_runs``, # which is defaulted as ``max_active_runs_per_dag``. -max_active_runs_per_dag = 16 +# max_active_runs_per_dag = 16 # The name of the method used in order to start Python processes via the multiprocessing module. # This corresponds directly with the options available in the Python docs: @@ -86,147 +86,147 @@ max_active_runs_per_dag = 16 # Whether to load the DAG examples that ship with Airflow. It's good to # get started, but you probably want to set this to ``False`` in a production # environment -load_examples = True +# load_examples = True # Path to the folder containing Airflow plugins -plugins_folder = {AIRFLOW_HOME}/plugins +# plugins_folder = {AIRFLOW_HOME}/plugins # Should tasks be executed via forking of the parent process ("False", # the speedier option) or by spawning a new python process ("True" slow, # but means plugin changes picked up by tasks straight away) -execute_tasks_new_python_interpreter = False +# execute_tasks_new_python_interpreter = False # Secret key to save connection passwords in the db -fernet_key = {FERNET_KEY} +# fernet_key = {FERNET_KEY} # Whether to disable pickling dags -donot_pickle = True +# donot_pickle = True # How long before timing out a python file import -dagbag_import_timeout = 30.0 +# dagbag_import_timeout = 30.0 # Should a traceback be shown in the UI for dagbag import errors, # instead of just the exception message -dagbag_import_error_tracebacks = True +# dagbag_import_error_tracebacks = True # If tracebacks are shown, how many entries from the traceback should be shown -dagbag_import_error_traceback_depth = 2 +# dagbag_import_error_traceback_depth = 2 # How long before timing out a DagFileProcessor, which processes a dag file -dag_file_processor_timeout = 50 +# dag_file_processor_timeout = 50 # The class to use for running task instances in a subprocess. # Choices include StandardTaskRunner, CgroupTaskRunner or the full import path to the class # when using a custom task runner. -task_runner = StandardTaskRunner +# task_runner = StandardTaskRunner # If set, tasks without a ``run_as_user`` argument will be run with this user # Can be used to de-elevate a sudo user running Airflow when executing tasks -default_impersonation = +# default_impersonation = # What security module to use (for example kerberos) -security = +# security = # Turn unit test mode on (overwrites many configuration options with test # values at runtime) -unit_test_mode = False +# unit_test_mode = False # Whether to enable pickling for xcom (note that this is insecure and allows for # RCE exploits). -enable_xcom_pickling = False +# enable_xcom_pickling = False # What classes can be imported during deserialization. This is a multi line value. # The individual items will be parsed as regexp. Python built-in classes (like dict) # are always allowed -allowed_deserialization_classes = airflow\..* +# allowed_deserialization_classes = airflow\..* # When a task is killed forcefully, this is the amount of time in seconds that # it has to cleanup after it is sent a SIGTERM, before it is SIGKILLED -killed_task_cleanup_time = 60 +# killed_task_cleanup_time = 60 # Whether to override params with dag_run.conf. If you pass some key-value pairs # through ``airflow dags backfill -c`` or # ``airflow dags trigger -c``, the key-value pairs will override the existing ones in params. -dag_run_conf_overrides_params = True +# dag_run_conf_overrides_params = True # When discovering DAGs, ignore any files that don't contain the strings ``DAG`` and ``airflow``. -dag_discovery_safe_mode = True +# dag_discovery_safe_mode = True # The pattern syntax used in the ".airflowignore" files in the DAG directories. Valid values are # ``regexp`` or ``glob``. -dag_ignore_file_syntax = regexp +# dag_ignore_file_syntax = regexp # The number of retries each task is going to have by default. Can be overridden at dag or task level. -default_task_retries = 0 +# default_task_retries = 0 # The number of seconds each task is going to wait by default between retries. Can be overridden at # dag or task level. -default_task_retry_delay = 300 +# default_task_retry_delay = 300 # The maximum delay (in seconds) each task is going to wait by default between retries. # This is a global setting and cannot be overridden at task or DAG level. -max_task_retry_delay = 86400 +# max_task_retry_delay = 86400 # The weighting method used for the effective total priority weight of the task -default_task_weight_rule = downstream +# default_task_weight_rule = downstream # The default task execution_timeout value for the operators. Expected an integer value to # be passed into timedelta as seconds. If not specified, then the value is considered as None, # meaning that the operators are never timed out by default. -default_task_execution_timeout = +# default_task_execution_timeout = # Updating serialized DAG can not be faster than a minimum interval to reduce database write rate. -min_serialized_dag_update_interval = 30 +# min_serialized_dag_update_interval = 30 # If True, serialized DAGs are compressed before writing to DB. # Note: this will disable the DAG dependencies view -compress_serialized_dags = False +# compress_serialized_dags = False # Fetching serialized DAG can not be faster than a minimum interval to reduce database # read rate. This config controls when your DAGs are updated in the Webserver -min_serialized_dag_fetch_interval = 10 +# min_serialized_dag_fetch_interval = 10 # Maximum number of Rendered Task Instance Fields (Template Fields) per task to store # in the Database. # All the template_fields for each of Task Instance are stored in the Database. # Keeping this number small may cause an error when you try to view ``Rendered`` tab in # TaskInstance view for older tasks. -max_num_rendered_ti_fields_per_task = 30 +# max_num_rendered_ti_fields_per_task = 30 # On each dagrun check against defined SLAs -check_slas = True +# check_slas = True # Path to custom XCom class that will be used to store and resolve operators results # Example: xcom_backend = path.to.CustomXCom -xcom_backend = airflow.models.xcom.BaseXCom +# xcom_backend = airflow.models.xcom.BaseXCom # By default Airflow plugins are lazily-loaded (only loaded when required). Set it to ``False``, # if you want to load plugins whenever 'airflow' is invoked via cli or loaded from module. -lazy_load_plugins = True +# lazy_load_plugins = True # By default Airflow providers are lazily-discovered (discovery and imports happen only when required). # Set it to False, if you want to discover providers whenever 'airflow' is invoked via cli or # loaded from module. -lazy_discover_providers = True +# lazy_discover_providers = True # Hide sensitive Variables or Connection extra json keys from UI and task logs when set to True # # (Connection passwords are always hidden in logs) -hide_sensitive_var_conn_fields = True +# hide_sensitive_var_conn_fields = True # A comma-separated list of extra sensitive keywords to look for in variables names or connection's # extra JSON. -sensitive_var_conn_names = +# sensitive_var_conn_names = # Task Slot counts for ``default_pool``. This setting would not have any effect in an existing # deployment where the ``default_pool`` is already created. For existing deployments, users can # change the number of slots using Webserver, API or the CLI -default_pool_task_slot_count = 128 +# default_pool_task_slot_count = 128 # The maximum list/dict length an XCom can push to trigger task mapping. If the pushed list/dict has a # length exceeding this value, the task pushing the XCom will be failed automatically to prevent the # mapped tasks from clogging the scheduler. -max_map_length = 1024 +# max_map_length = 1024 # The default umask to use for process when run in daemon mode (scheduler, worker, etc.) # @@ -234,7 +234,7 @@ max_map_length = 1024 # for newly created files. # # This value is treated as an octal-integer. -daemon_umask = 0o077 +# daemon_umask = 0o077 # Class to use as dataset manager. # Example: dataset_manager_class = airflow.datasets.manager.DatasetManager @@ -245,7 +245,7 @@ daemon_umask = 0o077 # dataset_manager_kwargs = # (experimental) Whether components should use Airflow Internal API for DB connectivity. -database_access_isolation = False +# database_access_isolation = False # (experimental)Airflow Internal API url. Only used if [core] database_access_isolation is True. # Example: internal_api_url = http://localhost:8080 @@ -256,14 +256,14 @@ database_access_isolation = False # SqlAlchemy supports many different database engines. # More information here: # http://airflow.apache.org/docs/apache-airflow/stable/howto/set-up-database.html#database-uri -sql_alchemy_conn = sqlite:///{AIRFLOW_HOME}/airflow.db +# sql_alchemy_conn = sqlite:///{AIRFLOW_HOME}/airflow.db # Extra engine specific keyword args passed to SQLAlchemy's create_engine, as a JSON-encoded value # Example: sql_alchemy_engine_args = {{"arg1": True}} # sql_alchemy_engine_args = # The encoding for the databases -sql_engine_encoding = utf-8 +# sql_engine_encoding = utf-8 # Collation for ``dag_id``, ``task_id``, ``key``, ``external_executor_id`` columns # in case they have different encoding. @@ -274,11 +274,11 @@ sql_engine_encoding = utf-8 # sql_engine_collation_for_ids = # If SqlAlchemy should pool database connections. -sql_alchemy_pool_enabled = True +# sql_alchemy_pool_enabled = True # The SqlAlchemy pool size is the maximum number of database connections # in the pool. 0 indicates no limit. -sql_alchemy_pool_size = 5 +# sql_alchemy_pool_size = 5 # The maximum overflow size of the pool. # When the number of checked-out connections reaches the size set in pool_size, @@ -289,23 +289,23 @@ sql_alchemy_pool_size = 5 # and the total number of "sleeping" connections the pool will allow is pool_size. # max_overflow can be set to ``-1`` to indicate no overflow limit; # no limit will be placed on the total number of concurrent connections. Defaults to ``10``. -sql_alchemy_max_overflow = 10 +# sql_alchemy_max_overflow = 10 # The SqlAlchemy pool recycle is the number of seconds a connection # can be idle in the pool before it is invalidated. This config does # not apply to sqlite. If the number of DB connections is ever exceeded, # a lower config value will allow the system to recover faster. -sql_alchemy_pool_recycle = 1800 +# sql_alchemy_pool_recycle = 1800 # Check connection at the start of each connection pool checkout. # Typically, this is a simple statement like "SELECT 1". # More information here: # https://docs.sqlalchemy.org/en/14/core/pooling.html#disconnect-handling-pessimistic -sql_alchemy_pool_pre_ping = True +# sql_alchemy_pool_pre_ping = True # The schema to use for the metadata database. # SqlAlchemy supports databases with the concept of multiple schemas. -sql_alchemy_schema = +# sql_alchemy_schema = # Import path for connect args in SqlAlchemy. Defaults to an empty dict. # This is useful when you want to configure db engine args that SqlAlchemy won't parse @@ -316,12 +316,12 @@ sql_alchemy_schema = # Whether to load the default connections that ship with Airflow. It's good to # get started, but you probably want to set this to ``False`` in a production # environment -load_default_connections = True +# load_default_connections = True # Number of times the code should be retried in case of DB Operational Errors. # Not all transactions will be retried as it can cause undesired state. # Currently it is only used in ``DagFileProcessor.process_file`` to retry ``dagbag.sync_to_db``. -max_db_retries = 3 +# max_db_retries = 3 [logging] # The folder where airflow should store its log files. @@ -329,22 +329,22 @@ max_db_retries = 3 # There are a few existing configurations that assume this is set to the default. # If you choose to override this you may need to update the dag_processor_manager_log_location and # dag_processor_manager_log_location settings as well. -base_log_folder = {AIRFLOW_HOME}/logs +# base_log_folder = {AIRFLOW_HOME}/logs # Airflow can store logs remotely in AWS S3, Google Cloud Storage or Elastic Search. # Set this to True if you want to enable remote logging. -remote_logging = False +# remote_logging = False # Users must supply an Airflow connection id that provides access to the storage # location. Depending on your remote logging service, this may only be used for # reading logs, not writing them. -remote_log_conn_id = +# remote_log_conn_id = # Path to Google Credential JSON file. If omitted, authorization based on `the Application Default # Credentials # `__ will # be used. -google_key_path = +# google_key_path = # Storage bucket URL for remote logging # S3 buckets should start with "s3://" @@ -352,50 +352,50 @@ google_key_path = # GCS buckets should start with "gs://" # WASB buckets should start with "wasb" just to help Airflow select correct handler # Stackdriver logs should start with "stackdriver://" -remote_base_log_folder = +# remote_base_log_folder = # Use server-side encryption for logs stored in S3 -encrypt_s3_logs = False +# encrypt_s3_logs = False # Logging level. # # Supported values: ``CRITICAL``, ``ERROR``, ``WARNING``, ``INFO``, ``DEBUG``. -logging_level = INFO +# logging_level = INFO # Logging level for celery. If not set, it uses the value of logging_level # # Supported values: ``CRITICAL``, ``ERROR``, ``WARNING``, ``INFO``, ``DEBUG``. -celery_logging_level = +# celery_logging_level = # Logging level for Flask-appbuilder UI. # # Supported values: ``CRITICAL``, ``ERROR``, ``WARNING``, ``INFO``, ``DEBUG``. -fab_logging_level = WARNING +# fab_logging_level = WARNING # Logging class # Specify the class that will specify the logging configuration # This class has to be on the python classpath # Example: logging_config_class = my.path.default_local_settings.LOGGING_CONFIG -logging_config_class = +# logging_config_class = # Flag to enable/disable Colored logs in Console # Colour the logs when the controlling terminal is a TTY. -colored_console_log = True +# colored_console_log = True # Log format for when Colored logs is enabled -colored_log_format = [%%(blue)s%%(asctime)s%%(reset)s] {{%%(blue)s%%(filename)s:%%(reset)s%%(lineno)d}} %%(log_color)s%%(levelname)s%%(reset)s - %%(log_color)s%%(message)s%%(reset)s -colored_formatter_class = airflow.utils.log.colored_log.CustomTTYColoredFormatter +# colored_log_format = [%%(blue)s%%(asctime)s%%(reset)s] {{%%(blue)s%%(filename)s:%%(reset)s%%(lineno)d}} %%(log_color)s%%(levelname)s%%(reset)s - %%(log_color)s%%(message)s%%(reset)s +# colored_formatter_class = airflow.utils.log.colored_log.CustomTTYColoredFormatter # Format of Log line -log_format = [%%(asctime)s] {{%%(filename)s:%%(lineno)d}} %%(levelname)s - %%(message)s -simple_log_format = %%(asctime)s %%(levelname)s - %%(message)s +# log_format = [%%(asctime)s] {{%%(filename)s:%%(lineno)d}} %%(levelname)s - %%(message)s +# simple_log_format = %%(asctime)s %%(levelname)s - %%(message)s # Where to send dag parser logs. If "file", logs are sent to log files defined by child_process_log_directory. -dag_processor_log_target = file +# dag_processor_log_target = file # Format of Dag Processor Log line -dag_processor_log_format = [%%(asctime)s] [SOURCE:DAG_PROCESSOR] {{%%(filename)s:%%(lineno)d}} %%(levelname)s - %%(message)s -log_formatter_class = airflow.utils.log.timezone_aware.TimezoneAware +# dag_processor_log_format = [%%(asctime)s] [SOURCE:DAG_PROCESSOR] {{%%(filename)s:%%(lineno)d}} %%(levelname)s - %%(message)s +# log_formatter_class = airflow.utils.log.timezone_aware.TimezoneAware # An import path to a function to add adaptations of each secret added with # `airflow.utils.log.secrets_masker.mask_secret` to be masked in log messages. The given function @@ -403,63 +403,63 @@ log_formatter_class = airflow.utils.log.timezone_aware.TimezoneAware # single adaptation of the secret or an iterable of adaptations to each be masked as secrets. # The original secret will be masked as well as any adaptations returned. # Example: secret_mask_adapter = urllib.parse.quote -secret_mask_adapter = +# secret_mask_adapter = # Specify prefix pattern like mentioned below with stream handler TaskHandlerWithCustomFormatter # Example: task_log_prefix_template = {{ti.dag_id}}-{{ti.task_id}}-{{execution_date}}-{{try_number}} -task_log_prefix_template = +# task_log_prefix_template = # Formatting for how airflow generates file names/paths for each task run. -log_filename_template = dag_id={{{{ ti.dag_id }}}}/run_id={{{{ ti.run_id }}}}/task_id={{{{ ti.task_id }}}}/{{%% if ti.map_index >= 0 %%}}map_index={{{{ ti.map_index }}}}/{{%% endif %%}}attempt={{{{ try_number }}}}.log +# log_filename_template = dag_id={{{{ ti.dag_id }}}}/run_id={{{{ ti.run_id }}}}/task_id={{{{ ti.task_id }}}}/{{%% if ti.map_index >= 0 %%}}map_index={{{{ ti.map_index }}}}/{{%% endif %%}}attempt={{{{ try_number }}}}.log # Formatting for how airflow generates file names for log -log_processor_filename_template = {{{{ filename }}}}.log +# log_processor_filename_template = {{{{ filename }}}}.log # Full path of dag_processor_manager logfile. -dag_processor_manager_log_location = {AIRFLOW_HOME}/logs/dag_processor_manager/dag_processor_manager.log +# dag_processor_manager_log_location = {AIRFLOW_HOME}/logs/dag_processor_manager/dag_processor_manager.log # Name of handler to read task instance logs. # Defaults to use ``task`` handler. -task_log_reader = task +# task_log_reader = task # A comma\-separated list of third-party logger names that will be configured to print messages to # consoles\. # Example: extra_logger_names = connexion,sqlalchemy -extra_logger_names = +# extra_logger_names = # When you start an airflow worker, airflow starts a tiny web server # subprocess to serve the workers local log files to the airflow main # web server, who then builds pages and sends them to users. This defines # the port on which the logs are served. It needs to be unused, and open # visible from the main web server to connect into the workers. -worker_log_server_port = 8793 +# worker_log_server_port = 8793 [metrics] # StatsD (https://github.com/etsy/statsd) integration settings. # Enables sending metrics to StatsD. -statsd_on = False -statsd_host = localhost -statsd_port = 8125 -statsd_prefix = airflow +# statsd_on = False +# statsd_host = localhost +# statsd_port = 8125 +# statsd_prefix = airflow # If you want to avoid sending all the available metrics to StatsD, # you can configure an allow list of prefixes (comma separated) to send only the metrics that # start with the elements of the list (e.g: "scheduler,executor,dagrun") -statsd_allow_list = +# statsd_allow_list = # A function that validate the StatsD stat name, apply changes to the stat name if necessary and return # the transformed stat name. # # The function should have the following signature: # def func_name(stat_name: str) -> str: -stat_name_handler = +# stat_name_handler = # To enable datadog integration to send airflow metrics. -statsd_datadog_enabled = False +# statsd_datadog_enabled = False # List of datadog tags attached to all metrics(e.g: key1:value1,key2:value2) -statsd_datadog_tags = +# statsd_datadog_tags = # If you want to utilise your own custom StatsD client set the relevant # module path below. @@ -469,29 +469,29 @@ statsd_datadog_tags = [secrets] # Full class name of secrets backend to enable (will precede env vars and metastore in search path) # Example: backend = airflow.providers.amazon.aws.secrets.systems_manager.SystemsManagerParameterStoreBackend -backend = +# backend = # The backend_kwargs param is loaded into a dictionary and passed to __init__ of secrets backend class. # See documentation for the secrets backend you are using. JSON is expected. # Example for AWS Systems Manager ParameterStore: # ``{{"connections_prefix": "/airflow/connections", "profile_name": "default"}}`` -backend_kwargs = +# backend_kwargs = [cli] # In what way should the cli access the API. The LocalClient will use the # database directly, while the json_client will use the api running on the # webserver -api_client = airflow.api.client.local_client +# api_client = airflow.api.client.local_client # If you set web_server_url_prefix, do NOT forget to append it here, ex: # ``endpoint_url = http://localhost:8080/myroot`` # So api will look like: ``http://localhost:8080/myroot/api/experimental/...`` -endpoint_url = http://localhost:8080 +# endpoint_url = http://localhost:8080 [debug] # Used only with ``DebugExecutor``. If set to ``True`` DAG will fail with first # failed task. Helpful for debugging purposes. -fail_fast = False +# fail_fast = False [api] # Enables the deprecated experimental API. Please note that these APIs do not have access control. @@ -504,76 +504,76 @@ fail_fast = False # `the Stable REST API `__. # For more information on migration, see # `RELEASE_NOTES.rst `_ -enable_experimental_api = False +# enable_experimental_api = False # Comma separated list of auth backends to authenticate users of the API. See # https://airflow.apache.org/docs/apache-airflow/stable/security/api.html for possible values. # ("airflow.api.auth.backend.default" allows all requests for historic reasons) -auth_backends = airflow.api.auth.backend.session +# auth_backends = airflow.api.auth.backend.session # Used to set the maximum page limit for API requests -maximum_page_limit = 100 +# maximum_page_limit = 100 # Used to set the default page limit when limit is zero. A default limit # of 100 is set on OpenApi spec. However, this particular default limit # only work when limit is set equal to zero(0) from API requests. # If no limit is supplied, the OpenApi spec default is used. -fallback_page_limit = 100 +# fallback_page_limit = 100 # The intended audience for JWT token credentials used for authorization. This value must match on the client and server sides. If empty, audience will not be tested. # Example: google_oauth2_audience = project-id-random-value.apps.googleusercontent.com -google_oauth2_audience = +# google_oauth2_audience = # Path to Google Cloud Service Account key file (JSON). If omitted, authorization based on # `the Application Default Credentials # `__ will # be used. # Example: google_key_path = /files/service-account-json -google_key_path = +# google_key_path = # Used in response to a preflight request to indicate which HTTP # headers can be used when making the actual request. This header is # the server side response to the browser's # Access-Control-Request-Headers header. -access_control_allow_headers = +# access_control_allow_headers = # Specifies the method or methods allowed when accessing the resource. -access_control_allow_methods = +# access_control_allow_methods = # Indicates whether the response can be shared with requesting code from the given origins. # Separate URLs with space. -access_control_allow_origins = +# access_control_allow_origins = [lineage] # what lineage backend to use -backend = +# backend = [atlas] -sasl_enabled = False -host = -port = 21000 -username = -password = +# sasl_enabled = False +# host = +# port = 21000 +# username = +# password = [operators] # The default owner assigned to each new operator, unless # provided explicitly or passed via ``default_args`` -default_owner = airflow -default_cpus = 1 -default_ram = 512 -default_disk = 512 -default_gpus = 0 +# default_owner = airflow +# default_cpus = 1 +# default_ram = 512 +# default_disk = 512 +# default_gpus = 0 # Default queue that tasks get assigned to and that worker listen on. -default_queue = default +# default_queue = default # Is allowed to pass additional/unused arguments (args, kwargs) to the BaseOperator operator. # If set to False, an exception will be thrown, otherwise only the console message will be displayed. -allow_illegal_arguments = False +# allow_illegal_arguments = False [hive] # Default mapreduce queue for HiveOperator tasks -default_hive_mapred_queue = +# default_hive_mapred_queue = # Template for mapred_job_name in HiveOperator, supports the following named parameters # hostname, dag_id, task_id, execution_date @@ -583,49 +583,49 @@ default_hive_mapred_queue = # The base url of your website as airflow cannot guess what domain or # cname you are using. This is used in automated emails that # airflow sends to point links to the right web server -base_url = http://localhost:8080 +# base_url = http://localhost:8080 # Default timezone to display all dates in the UI, can be UTC, system, or # any IANA timezone string (e.g. Europe/Amsterdam). If left empty the # default value of core/default_timezone will be used # Example: default_ui_timezone = America/New_York -default_ui_timezone = UTC +# default_ui_timezone = UTC # The ip specified when starting the web server -web_server_host = 0.0.0.0 +# web_server_host = 0.0.0.0 # The port on which to run the web server -web_server_port = 8080 +# web_server_port = 8080 # Paths to the SSL certificate and key for the web server. When both are # provided SSL will be enabled. This does not change the web server port. -web_server_ssl_cert = +# web_server_ssl_cert = # Paths to the SSL certificate and key for the web server. When both are # provided SSL will be enabled. This does not change the web server port. -web_server_ssl_key = +# web_server_ssl_key = # The type of backend used to store web session data, can be 'database' or 'securecookie' # Example: session_backend = securecookie -session_backend = database +# session_backend = database # Number of seconds the webserver waits before killing gunicorn master that doesn't respond -web_server_master_timeout = 120 +# web_server_master_timeout = 120 # Number of seconds the gunicorn webserver waits before timing out on a worker -web_server_worker_timeout = 120 +# web_server_worker_timeout = 120 # Number of workers to refresh at a time. When set to 0, worker refresh is # disabled. When nonzero, airflow periodically refreshes webserver workers by # bringing up new ones and killing old ones. -worker_refresh_batch_size = 1 +# worker_refresh_batch_size = 1 # Number of seconds to wait before refreshing a batch of workers. -worker_refresh_interval = 6000 +# worker_refresh_interval = 6000 # If set to True, Airflow will track files in plugins_folder directory. When it detects changes, # then reload the gunicorn. -reload_on_plugin_change = False +# reload_on_plugin_change = False # Secret key used to run your flask app. It should be as random as possible. However, when running # more than 1 instances of webserver, make sure all of them use the same ``secret_key`` otherwise @@ -634,100 +634,100 @@ reload_on_plugin_change = False # The token generated using the secret key has a short expiry time though - make sure that time on # ALL the machines that you run airflow components on is synchronized (for example using ntpd) # otherwise you might get "forbidden" errors when the logs are accessed. -secret_key = {SECRET_KEY} +# secret_key = {SECRET_KEY} # Number of workers to run the Gunicorn web server -workers = 4 +# workers = 4 # The worker class gunicorn should use. Choices include # sync (default), eventlet, gevent -worker_class = sync +# worker_class = sync # Log files for the gunicorn webserver. '-' means log to stderr. -access_logfile = - +# access_logfile = - # Log files for the gunicorn webserver. '-' means log to stderr. -error_logfile = - +# error_logfile = - # Access log format for gunicorn webserver. # default format is %%(h)s %%(l)s %%(u)s %%(t)s "%%(r)s" %%(s)s %%(b)s "%%(f)s" "%%(a)s" # documentation - https://docs.gunicorn.org/en/stable/settings.html#access-log-format -access_logformat = +# access_logformat = # Expose the configuration file in the web server. Set to "non-sensitive-only" to show all values # except those that have security implications. "True" shows all values. "False" hides the # configuration completely. -expose_config = False +# expose_config = False # Expose hostname in the web server -expose_hostname = True +# expose_hostname = True # Expose stacktrace in the web server -expose_stacktrace = False +# expose_stacktrace = False # Default DAG view. Valid values are: ``grid``, ``graph``, ``duration``, ``gantt``, ``landing_times`` -dag_default_view = grid +# dag_default_view = grid # Default DAG orientation. Valid values are: # ``LR`` (Left->Right), ``TB`` (Top->Bottom), ``RL`` (Right->Left), ``BT`` (Bottom->Top) -dag_orientation = LR +# dag_orientation = LR # The amount of time (in secs) webserver will wait for initial handshake # while fetching logs from other worker machine -log_fetch_timeout_sec = 5 +# log_fetch_timeout_sec = 5 # Time interval (in secs) to wait before next log fetching. -log_fetch_delay_sec = 2 +# log_fetch_delay_sec = 2 # Distance away from page bottom to enable auto tailing. -log_auto_tailing_offset = 30 +# log_auto_tailing_offset = 30 # Animation speed for auto tailing log display. -log_animation_speed = 1000 +# log_animation_speed = 1000 # By default, the webserver shows paused DAGs. Flip this to hide paused # DAGs by default -hide_paused_dags_by_default = False +# hide_paused_dags_by_default = False # Consistent page size across all listing views in the UI -page_size = 100 +# page_size = 100 # Define the color of navigation bar -navbar_color = #fff +# navbar_color = #fff # Default dagrun to show in UI -default_dag_run_display_number = 25 +# default_dag_run_display_number = 25 # Enable werkzeug ``ProxyFix`` middleware for reverse proxy -enable_proxy_fix = False +# enable_proxy_fix = False # Number of values to trust for ``X-Forwarded-For``. # More info: https://werkzeug.palletsprojects.com/en/0.16.x/middleware/proxy_fix/ -proxy_fix_x_for = 1 +# proxy_fix_x_for = 1 # Number of values to trust for ``X-Forwarded-Proto`` -proxy_fix_x_proto = 1 +# proxy_fix_x_proto = 1 # Number of values to trust for ``X-Forwarded-Host`` -proxy_fix_x_host = 1 +# proxy_fix_x_host = 1 # Number of values to trust for ``X-Forwarded-Port`` -proxy_fix_x_port = 1 +# proxy_fix_x_port = 1 # Number of values to trust for ``X-Forwarded-Prefix`` -proxy_fix_x_prefix = 1 +# proxy_fix_x_prefix = 1 # Set secure flag on session cookie -cookie_secure = False +# cookie_secure = False # Set samesite policy on session cookie -cookie_samesite = Lax +# cookie_samesite = Lax # Default setting for wrap toggle on DAG code and TI log views. -default_wrap = False +# default_wrap = False # Allow the UI to be rendered in a frame -x_frame_enabled = True +# x_frame_enabled = True # Send anonymous user activity to your analytics tool # choose from google_analytics, segment, or metarouter @@ -737,33 +737,33 @@ x_frame_enabled = True # analytics_id = # 'Recent Tasks' stats will show for old DagRuns if set -show_recent_stats_for_completed_runs = True +# show_recent_stats_for_completed_runs = True # Update FAB permissions and sync security manager roles # on webserver startup -update_fab_perms = True +# update_fab_perms = True # The UI cookie lifetime in minutes. User will be logged out from UI after # ``session_lifetime_minutes`` of non-activity -session_lifetime_minutes = 43200 +# session_lifetime_minutes = 43200 # Sets a custom page title for the DAGs overview page and site title for all pages # instance_name = # Whether the custom page title for the DAGs overview page contains any Markup language -instance_name_has_markup = False +# instance_name_has_markup = False # How frequently, in seconds, the DAG data will auto-refresh in graph or grid view # when auto-refresh is turned on -auto_refresh_interval = 3 +# auto_refresh_interval = 3 # Boolean for displaying warning for publicly viewable deployment -warn_deployment_exposure = True +# warn_deployment_exposure = True # Comma separated string of view events to exclude from dag audit view. # All other events will be added minus the ones passed here. # The audit logs in the db will not be affected by this parameter. -audit_view_excluded_events = gantt,landing_times,tries,duration,calendar,graph,grid,tree,tree_data +# audit_view_excluded_events = gantt,landing_times,tries,duration,calendar,graph,grid,tree,tree_data # Comma separated string of view events to include in dag audit view. # If passed, only these events will populate the dag audit view. @@ -772,26 +772,26 @@ audit_view_excluded_events = gantt,landing_times,tries,duration,calendar,graph,g # audit_view_included_events = # Boolean for running SwaggerUI in the webserver. -enable_swagger_ui = True +# enable_swagger_ui = True # Boolean for running Internal API in the webserver. -run_internal_api = False +# run_internal_api = False [email] # Configuration email backend and whether to # send email alerts on retry or failure # Email backend to use -email_backend = airflow.utils.email.send_email_smtp +# email_backend = airflow.utils.email.send_email_smtp # Email connection to use -email_conn_id = smtp_default +# email_conn_id = smtp_default # Whether email alerts should be sent when a task is retried -default_email_on_retry = True +# default_email_on_retry = True # Whether email alerts should be sent when a task failed -default_email_on_failure = True +# default_email_on_failure = True # File that will be used as the template for Email subject (which will be rendered using Jinja2). # If not set, Airflow uses a base template. @@ -813,17 +813,17 @@ default_email_on_failure = True # If you want airflow to send emails on retries, failure, and you want to use # the airflow.utils.email.send_email_smtp function, you have to configure an # smtp server here -smtp_host = localhost -smtp_starttls = True -smtp_ssl = False +# smtp_host = localhost +# smtp_starttls = True +# smtp_ssl = False # Example: smtp_user = airflow # smtp_user = # Example: smtp_password = airflow # smtp_password = -smtp_port = 25 -smtp_mail_from = airflow@example.com -smtp_timeout = 30 -smtp_retry_limit = 5 +# smtp_port = 25 +# smtp_mail_from = airflow@example.com +# smtp_timeout = 30 +# smtp_retry_limit = 5 [sentry] @@ -833,8 +833,8 @@ smtp_retry_limit = 5 # Unsupported options: ``integrations``, ``in_app_include``, ``in_app_exclude``, # ``ignore_errors``, ``before_breadcrumb``, ``transport``. # Enable error reporting to Sentry -sentry_on = false -sentry_dsn = +# sentry_on = false +# sentry_dsn = # Dotted path to a before_send function that the sentry SDK should be configured to use. # before_send = @@ -847,7 +847,7 @@ sentry_dsn = # When the queue of a task is the value of ``kubernetes_queue`` (default ``kubernetes``), # the task is executed via ``KubernetesExecutor``, # otherwise via ``LocalExecutor`` -kubernetes_queue = kubernetes +# kubernetes_queue = kubernetes [celery_kubernetes_executor] @@ -857,20 +857,20 @@ kubernetes_queue = kubernetes # When the queue of a task is the value of ``kubernetes_queue`` (default ``kubernetes``), # the task is executed via ``KubernetesExecutor``, # otherwise via ``CeleryExecutor`` -kubernetes_queue = kubernetes +# kubernetes_queue = kubernetes [celery] # This section only applies if you are using the CeleryExecutor in # ``[core]`` section above # The app name that will be used by celery -celery_app_name = airflow.executors.celery_executor +# celery_app_name = airflow.executors.celery_executor # The concurrency that will be used when starting workers with the # ``airflow celery worker`` command. This defines the number of task instances that # a worker will take, so size up your workers based on the resources on # your worker box and the nature of your tasks -worker_concurrency = 16 +# worker_concurrency = 16 # The maximum and minimum concurrency that will be used when starting workers with the # ``airflow celery worker`` command (always keep minimum processes, but grow @@ -888,16 +888,16 @@ worker_concurrency = 16 # running tasks while another worker has unutilized processes that are unable to process the already # claimed blocked tasks. # https://docs.celeryproject.org/en/stable/userguide/optimizing.html#prefetch-limits -worker_prefetch_multiplier = 1 +# worker_prefetch_multiplier = 1 # Specify if remote control of the workers is enabled. # When using Amazon SQS as the broker, Celery creates lots of ``.*reply-celery-pidbox`` queues. You can # prevent this by setting this to false. However, with this disabled Flower won't work. -worker_enable_remote_control = true +# worker_enable_remote_control = true # The Celery broker URL. Celery supports RabbitMQ, Redis and experimentally # a sqlalchemy database. Refer to the Celery documentation for more information. -broker_url = redis://redis:6379/0 +# broker_url = redis://redis:6379/0 # The Celery result_backend. When a job finishes, it needs to update the # metadata of the job. Therefore it will post a message on a message bus, @@ -911,64 +911,64 @@ broker_url = redis://redis:6379/0 # Celery Flower is a sweet UI for Celery. Airflow has a shortcut to start # it ``airflow celery flower``. This defines the IP that Celery Flower runs on -flower_host = 0.0.0.0 +# flower_host = 0.0.0.0 # The root URL for Flower # Example: flower_url_prefix = /flower -flower_url_prefix = +# flower_url_prefix = # This defines the port that Celery Flower runs on -flower_port = 5555 +# flower_port = 5555 # Securing Flower with Basic Authentication # Accepts user:password pairs separated by a comma # Example: flower_basic_auth = user1:password1,user2:password2 -flower_basic_auth = +# flower_basic_auth = # How many processes CeleryExecutor uses to sync task state. # 0 means to use max(1, number of cores - 1) processes. -sync_parallelism = 0 +# sync_parallelism = 0 # Import path for celery configuration options -celery_config_options = airflow.config_templates.default_celery.DEFAULT_CELERY_CONFIG -ssl_active = False -ssl_key = -ssl_cert = -ssl_cacert = +# celery_config_options = airflow.config_templates.default_celery.DEFAULT_CELERY_CONFIG +# ssl_active = False +# ssl_key = +# ssl_cert = +# ssl_cacert = # Celery Pool implementation. # Choices include: ``prefork`` (default), ``eventlet``, ``gevent`` or ``solo``. # See: # https://docs.celeryproject.org/en/latest/userguide/workers.html#concurrency # https://docs.celeryproject.org/en/latest/userguide/concurrency/eventlet.html -pool = prefork +# pool = prefork # The number of seconds to wait before timing out ``send_task_to_executor`` or # ``fetch_celery_task_state`` operations. -operation_timeout = 1.0 +# operation_timeout = 1.0 # Celery task will report its status as 'started' when the task is executed by a worker. # This is used in Airflow to keep track of the running tasks and if a Scheduler is restarted # or run in HA mode, it can adopt the orphan tasks launched by previous SchedulerJob. -task_track_started = True +# task_track_started = True # Time in seconds after which adopted tasks which are queued in celery are assumed to be stalled, # and are automatically rescheduled. This setting does the same thing as ``stalled_task_timeout`` but # applies specifically to adopted tasks only. When set to 0, the ``stalled_task_timeout`` setting # also applies to adopted tasks. -task_adoption_timeout = 600 +# task_adoption_timeout = 600 # Time in seconds after which tasks queued in celery are assumed to be stalled, and are automatically # rescheduled. Adopted tasks will instead use the ``task_adoption_timeout`` setting if specified. # When set to 0, automatic clearing of stalled tasks is disabled. -stalled_task_timeout = 0 +# stalled_task_timeout = 0 # The Maximum number of retries for publishing task messages to the broker when failing # due to ``AirflowTaskTimeout`` error before giving up and marking Task as failed. -task_publish_max_retries = 3 +# task_publish_max_retries = 3 # Worker initialisation check to validate Metadata Database connection -worker_precheck = False +# worker_precheck = False [celery_broker_transport_options] @@ -990,76 +990,76 @@ worker_precheck = False # This section only applies if you are using the DaskExecutor in # [core] section above # The IP address and port of the Dask cluster's scheduler. -cluster_address = 127.0.0.1:8786 +# cluster_address = 127.0.0.1:8786 # TLS/ SSL settings to access a secured Dask scheduler. -tls_ca = -tls_cert = -tls_key = +# tls_ca = +# tls_cert = +# tls_key = [scheduler] # Task instances listen for external kill signal (when you clear tasks # from the CLI or the UI), this defines the frequency at which they should # listen (in seconds). -job_heartbeat_sec = 5 +# job_heartbeat_sec = 5 # The scheduler constantly tries to trigger new tasks (look at the # scheduler section in the docs for more information). This defines # how often the scheduler should run (in seconds). -scheduler_heartbeat_sec = 5 +# scheduler_heartbeat_sec = 5 # The number of times to try to schedule each DAG file # -1 indicates unlimited number -num_runs = -1 +# num_runs = -1 # Controls how long the scheduler will sleep between loops, but if there was nothing to do # in the loop. i.e. if it scheduled something then it will start the next loop # iteration straight away. -scheduler_idle_sleep_time = 1 +# scheduler_idle_sleep_time = 1 # Number of seconds after which a DAG file is parsed. The DAG file is parsed every # ``min_file_process_interval`` number of seconds. Updates to DAGs are reflected after # this interval. Keeping this number low will increase CPU usage. -min_file_process_interval = 30 +# min_file_process_interval = 30 # How often (in seconds) to check for stale DAGs (DAGs which are no longer present in # the expected files) which should be deactivated, as well as datasets that are no longer # referenced and should be marked as orphaned. -parsing_cleanup_interval = 60 +# parsing_cleanup_interval = 60 # How often (in seconds) to scan the DAGs directory for new files. Default to 5 minutes. -dag_dir_list_interval = 300 +# dag_dir_list_interval = 300 # How often should stats be printed to the logs. Setting to 0 will disable printing stats -print_stats_interval = 30 +# print_stats_interval = 30 # How often (in seconds) should pool usage stats be sent to StatsD (if statsd_on is enabled) -pool_metrics_interval = 5.0 +# pool_metrics_interval = 5.0 # If the last scheduler heartbeat happened more than scheduler_health_check_threshold # ago (in seconds), scheduler is considered unhealthy. # This is used by the health check in the "/health" endpoint -scheduler_health_check_threshold = 30 +# scheduler_health_check_threshold = 30 # When you start a scheduler, airflow starts a tiny web server # subprocess to serve a health check if this is set to True -enable_health_check = False +# enable_health_check = False # When you start a scheduler, airflow starts a tiny web server # subprocess to serve a health check on this port -scheduler_health_check_server_port = 8974 +# scheduler_health_check_server_port = 8974 # How often (in seconds) should the scheduler check for orphaned tasks and SchedulerJobs -orphaned_tasks_check_interval = 300.0 -child_process_log_directory = {AIRFLOW_HOME}/logs/scheduler +# orphaned_tasks_check_interval = 300.0 +# child_process_log_directory = {AIRFLOW_HOME}/logs/scheduler # Local task jobs periodically heartbeat to the DB. If the job has # not heartbeat in this many seconds, the scheduler will mark the # associated task instance as failed and will re-schedule the task. -scheduler_zombie_task_threshold = 300 +# scheduler_zombie_task_threshold = 300 # How often (in seconds) should the scheduler check for zombie tasks. -zombie_detection_interval = 10.0 +# zombie_detection_interval = 10.0 # Turn off scheduler catchup by setting this to ``False``. # Default behavior is unchanged and @@ -1067,42 +1067,42 @@ zombie_detection_interval = 10.0 # will not do scheduler catchup if this is ``False``, # however it can be set on a per DAG basis in the # DAG definition (catchup) -catchup_by_default = True +# catchup_by_default = True # Setting this to True will make first task instance of a task # ignore depends_on_past setting. A task instance will be considered # as the first task instance of a task when there is no task instance # in the DB with an execution_date earlier than it., i.e. no manual marking # success will be needed for a newly added task to be scheduled. -ignore_first_depends_on_past_by_default = True +# ignore_first_depends_on_past_by_default = True # This changes the batch size of queries in the scheduling main loop. # If this is too high, SQL query performance may be impacted by # complexity of query predicate, and/or excessive locking. # Additionally, you may hit the maximum allowable query length for your db. # Set this to 0 for no limit (not advised) -max_tis_per_query = 512 +# max_tis_per_query = 512 # Should the scheduler issue ``SELECT ... FOR UPDATE`` in relevant queries. # If this is set to False then you should not run more than a single # scheduler at once -use_row_level_locking = True +# use_row_level_locking = True # Max number of DAGs to create DagRuns for per scheduler loop. -max_dagruns_to_create_per_loop = 10 +# max_dagruns_to_create_per_loop = 10 # How many DagRuns should a scheduler examine (and lock) when scheduling # and queuing tasks. -max_dagruns_per_loop_to_schedule = 20 +# max_dagruns_per_loop_to_schedule = 20 # Should the Task supervisor process perform a "mini scheduler" to attempt to schedule more tasks of the # same DAG. Leaving this on will mean tasks in the same DAG execute quicker, but might starve out other # dags in some circumstances -schedule_after_task_execution = True +# schedule_after_task_execution = True # The scheduler can run multiple processes in parallel to parse dags. # This defines how many processes will run. -parsing_processes = 2 +# parsing_processes = 2 # One of ``modified_time``, ``random_seeded_by_host`` and ``alphabetical``. # The scheduler will list and sort the dag files to decide the parsing order. @@ -1113,132 +1113,132 @@ parsing_processes = 2 # same host. This is useful when running with Scheduler in HA mode where each scheduler can # parse different DAG files. # * ``alphabetical``: Sort by filename -file_parsing_sort_mode = modified_time +# file_parsing_sort_mode = modified_time # Whether the dag processor is running as a standalone process or it is a subprocess of a scheduler # job. -standalone_dag_processor = False +# standalone_dag_processor = False # Only applicable if `[scheduler]standalone_dag_processor` is true and callbacks are stored # in database. Contains maximum number of callbacks that are fetched during a single loop. -max_callbacks_per_loop = 20 +# max_callbacks_per_loop = 20 # Only applicable if `[scheduler]standalone_dag_processor` is true. # Time in seconds after which dags, which were not updated by Dag Processor are deactivated. -dag_stale_not_seen_duration = 600 +# dag_stale_not_seen_duration = 600 # Turn off scheduler use of cron intervals by setting this to False. # DAGs submitted manually in the web UI or with trigger_dag will still run. -use_job_schedule = True +# use_job_schedule = True # Allow externally triggered DagRuns for Execution Dates in the future # Only has effect if schedule_interval is set to None in DAG -allow_trigger_in_future = False +# allow_trigger_in_future = False # How often to check for expired trigger requests that have not run yet. -trigger_timeout_check_interval = 15 +# trigger_timeout_check_interval = 15 [triggerer] # How many triggers a single Triggerer will run at once, by default. -default_capacity = 1000 +# default_capacity = 1000 [kerberos] -ccache = /tmp/airflow_krb5_ccache +# ccache = /tmp/airflow_krb5_ccache # gets augmented with fqdn -principal = airflow -reinit_frequency = 3600 -kinit_path = kinit -keytab = airflow.keytab +# principal = airflow +# reinit_frequency = 3600 +# kinit_path = kinit +# keytab = airflow.keytab # Allow to disable ticket forwardability. -forwardable = True +# forwardable = True # Allow to remove source IP from token, useful when using token behind NATted Docker host. -include_ip = True +# include_ip = True [elasticsearch] # Elasticsearch host -host = +# host = # Format of the log_id, which is used to query for a given tasks logs -log_id_template = {{dag_id}}-{{task_id}}-{{run_id}}-{{map_index}}-{{try_number}} +# log_id_template = {{dag_id}}-{{task_id}}-{{run_id}}-{{map_index}}-{{try_number}} # Used to mark the end of a log stream for a task -end_of_log_mark = end_of_log +# end_of_log_mark = end_of_log # Qualified URL for an elasticsearch frontend (like Kibana) with a template argument for log_id # Code will construct log_id using the log_id template from the argument above. # NOTE: scheme will default to https if one is not provided # Example: frontend = http://localhost:5601/app/kibana#/discover?_a=(columns:!(message),query:(language:kuery,query:'log_id: "{{log_id}}"'),sort:!(log.offset,asc)) -frontend = +# frontend = # Write the task logs to the stdout of the worker, rather than the default files -write_stdout = False +# write_stdout = False # Instead of the default log formatter, write the log lines as JSON -json_format = False +# json_format = False # Log fields to also attach to the json output, if enabled -json_fields = asctime, filename, lineno, levelname, message +# json_fields = asctime, filename, lineno, levelname, message # The field where host name is stored (normally either `host` or `host.name`) -host_field = host +# host_field = host # The field where offset is stored (normally either `offset` or `log.offset`) -offset_field = offset +# offset_field = offset # Comma separated list of index patterns to use when searching for logs (default: `_all`). # Example: index_patterns = something-* -index_patterns = _all +# index_patterns = _all [elasticsearch_configs] -use_ssl = False -verify_certs = True +# use_ssl = False +# verify_certs = True [kubernetes_executor] # Path to the YAML pod file that forms the basis for KubernetesExecutor workers. -pod_template_file = +# pod_template_file = # The repository of the Kubernetes Image for the Worker to Run -worker_container_repository = +# worker_container_repository = # The tag of the Kubernetes Image for the Worker to Run -worker_container_tag = +# worker_container_tag = # The Kubernetes namespace where airflow workers should be created. Defaults to ``default`` -namespace = default +# namespace = default # If True, all worker pods will be deleted upon termination -delete_worker_pods = True +# delete_worker_pods = True # If False (and delete_worker_pods is True), # failed worker pods will not be deleted so users can investigate them. # This only prevents removal of worker pods where the worker itself failed, # not when the task it ran failed. -delete_worker_pods_on_failure = False +# delete_worker_pods_on_failure = False # Number of Kubernetes Worker Pod creation calls per scheduler loop. # Note that the current default of "1" will only launch a single pod # per-heartbeat. It is HIGHLY recommended that users increase this # number to match the tolerance of their kubernetes cluster for # better performance. -worker_pods_creation_batch_size = 1 +# worker_pods_creation_batch_size = 1 # Allows users to launch pods in multiple namespaces. # Will require creating a cluster-role for the scheduler, # or use multi_namespace_mode_namespace_list configuration. -multi_namespace_mode = False +# multi_namespace_mode = False # If multi_namespace_mode is True while scheduler does not have a cluster-role, # give the list of namespaces where the scheduler will schedule jobs # Scheduler needs to have the necessary permissions in these namespaces. -multi_namespace_mode_namespace_list = +# multi_namespace_mode_namespace_list = # Use the service account kubernetes gives to pods to connect to kubernetes cluster. # It's intended for clients that expect to be running inside a pod running on kubernetes. # It will raise an exception if called from a process not running in a kubernetes environment. -in_cluster = True +# in_cluster = True # When running with in_cluster=False change the default cluster_context or config_file # options to Kubernetes client. Leave blank these to use default behaviour like ``kubectl`` has. @@ -1252,7 +1252,7 @@ in_cluster = True # List of supported params are similar for all core_v1_apis, hence a single config # variable for all apis. See: # https://raw.githubusercontent.com/kubernetes-client/python/41f11a09995efcd0142e25946adc7591431bfb2f/kubernetes/client/api/core_v1_api.py -kube_client_request_args = +# kube_client_request_args = # Optional keyword arguments to pass to the ``delete_namespaced_pod`` kubernetes client # ``core_v1_api`` method when using the Kubernetes Executor. @@ -1260,41 +1260,41 @@ kube_client_request_args = # class defined here: # https://github.com/kubernetes-client/python/blob/41f11a09995efcd0142e25946adc7591431bfb2f/kubernetes/client/models/v1_delete_options.py#L19 # Example: delete_option_kwargs = {{"grace_period_seconds": 10}} -delete_option_kwargs = +# delete_option_kwargs = # Enables TCP keepalive mechanism. This prevents Kubernetes API requests to hang indefinitely # when idle connection is time-outed on services like cloud load balancers or firewalls. -enable_tcp_keepalive = True +# enable_tcp_keepalive = True # When the `enable_tcp_keepalive` option is enabled, TCP probes a connection that has # been idle for `tcp_keep_idle` seconds. -tcp_keep_idle = 120 +# tcp_keep_idle = 120 # When the `enable_tcp_keepalive` option is enabled, if Kubernetes API does not respond # to a keepalive probe, TCP retransmits the probe after `tcp_keep_intvl` seconds. -tcp_keep_intvl = 30 +# tcp_keep_intvl = 30 # When the `enable_tcp_keepalive` option is enabled, if Kubernetes API does not respond # to a keepalive probe, TCP retransmits the probe `tcp_keep_cnt number` of times before # a connection is considered to be broken. -tcp_keep_cnt = 6 +# tcp_keep_cnt = 6 # Set this to false to skip verifying SSL certificate of Kubernetes python client. -verify_ssl = True +# verify_ssl = True # How long in seconds a worker can be in Pending before it is considered a failure -worker_pods_pending_timeout = 300 +# worker_pods_pending_timeout = 300 # How often in seconds to check if Pending workers have exceeded their timeouts -worker_pods_pending_timeout_check_interval = 120 +# worker_pods_pending_timeout_check_interval = 120 # How often in seconds to check for task instances stuck in "queued" status without a pod -worker_pods_queued_check_interval = 60 +# worker_pods_queued_check_interval = 60 # How many pending pods to check for timeout violations in each check interval. # You may want this higher if you have a very large cluster and/or use ``multi_namespace_mode``. -worker_pods_pending_timeout_batch_size = 100 +# worker_pods_pending_timeout_batch_size = 100 [sensors] # Sensor default timeout, 7 days by default (7 * 24 * 60 * 60). -default_timeout = 604800 +# default_timeout = 604800 diff --git a/airflow/configuration.py b/airflow/configuration.py index 712f40c531fe0..a51a19ff55db1 100644 --- a/airflow/configuration.py +++ b/airflow/configuration.py @@ -42,6 +42,7 @@ from typing_extensions import overload from airflow.compat.functools import cached_property +from airflow.config_templates.built_in_defaults import built_in_defaults from airflow.exceptions import AirflowConfigException from airflow.secrets import DEFAULT_SECRETS_SEARCH_PATH, BaseSecretsBackend from airflow.utils import yaml @@ -309,27 +310,20 @@ def inversed_deprecated_sections(self): def optionxform(self, optionstr: str) -> str: return optionstr - def __init__(self, default_config: str | None = None, *args, **kwargs): + def __init__(self, *args, **kwargs): super().__init__(*args, **kwargs) self.upgraded_values = {} - - self.airflow_defaults = ConfigParser(*args, **kwargs) - if default_config is not None: - self.airflow_defaults.read_string(default_config) - # Set the upgrade value based on the current loaded default - default = self.airflow_defaults.get("logging", "log_filename_template", fallback=None) - if default: - replacement = self.deprecated_values["logging"]["log_filename_template"] - self.deprecated_values["logging"]["log_filename_template"] = ( - replacement[0], - default, - replacement[2], - ) - else: - # In case of tests it might not exist - with suppress(KeyError): - del self.deprecated_values["logging"]["log_filename_template"] + # Set the upgrade value based on the built-in defaults + default_template = built_in_defaults.get("logging").get("log_filename_template") + if default_template: + replacement = self.deprecated_values["logging"]["log_filename_template"] + self.deprecated_values["logging"]["log_filename_template"] = ( + replacement[0], + default_template, + replacement[2], + ) else: + # In case of tests it might not exist with suppress(KeyError): del self.deprecated_values["logging"]["log_filename_template"] @@ -343,7 +337,7 @@ def validate(self): for section, replacement in self.deprecated_values.items(): for name, info in replacement.items(): old, new, version = info - current_value = self.get(section, name, fallback="") + current_value = self.get(section, name) if self._using_old_value(old, current_value): self.upgraded_values[(section, name)] = current_value new_value = old.sub(new, current_value) @@ -629,12 +623,14 @@ def get(self, section: str, key: str, **kwargs) -> str | None: # type: ignore[o if option is not None: return option - # ...then the default config - if self.airflow_defaults.has_option(section, key) or "fallback" in kwargs: - return expand_env_var(self.airflow_defaults.get(section, key, **kwargs)) + # ... then built in defaults + section_dict = built_in_defaults.get(section) + if section_dict is not None: + val = section_dict.get(key) + if val is not None: + return expand_env_var(val) log.warning("section/key [%s/%s] not found in config", section, key) - raise AirflowConfigException(f"section/key [{section}/{key}] not found in config") def _get_option_from_secrets( @@ -882,8 +878,8 @@ def remove_option(self, section: str, option: str, remove_default: bool = True): if super().has_option(section, option): super().remove_option(section, option) - if self.airflow_defaults.has_option(section, option) and remove_default: - self.airflow_defaults.remove_option(section, option) + if built_in_defaults.get(section) and built_in_defaults.get(section).get(option) and remove_default: + del built_in_defaults.get(section)[option] def getsection(self, section: str) -> ConfigOptionsDictType | None: """ @@ -893,15 +889,14 @@ def getsection(self, section: str) -> ConfigOptionsDictType | None: :param section: section from the config """ - if not self.has_section(section) and not self.airflow_defaults.has_section(section): + if not self.has_section(section) and not built_in_defaults.get(section): return None - if self.airflow_defaults.has_section(section): - _section: ConfigOptionsDictType = OrderedDict(self.airflow_defaults.items(section)) - else: - _section = OrderedDict() + _section: dict = {} + if built_in_defaults.get(section): + _section.update(built_in_defaults.get(section)) if self.has_section(section): - _section.update(OrderedDict(self.items(section))) + _section.update(self.items(section)) section_prefix = self._env_var_name(section, "") for env_var in sorted(os.environ.keys()): @@ -940,10 +935,6 @@ def write( # type: ignore[override] delimiter = f" {self._delimiters[0]} " # type: ignore[attr-defined] else: delimiter = self._delimiters[0] # type: ignore[attr-defined] - if self._defaults: # type: ignore - self._write_section( # type: ignore[attr-defined] - fp, self.default_section, self._defaults.items(), delimiter # type: ignore[attr-defined] - ) sections = ( {section: dict(self.getsection(section))} # type: ignore[arg-type] if section @@ -1004,7 +995,7 @@ def as_dict( config_sources: ConfigSourcesType = {} configs = [ - ("default", self.airflow_defaults), + ("default", built_in_defaults), ("airflow.cfg", self), ] @@ -1171,7 +1162,7 @@ def _filter_by_source( if not getter_opt: continue # Check to see that there is a default value - if not self.airflow_defaults.has_option(section, key): + if not built_in_defaults.get(section) and built_in_defaults.get(section).get(key): continue # Check to see if bare setting is the same as defaults if display_source: @@ -1179,7 +1170,7 @@ def _filter_by_source( opt, source = config_sources[section][key] # type: ignore else: opt = config_sources[section][key] - if opt == self.airflow_defaults.get(section, key): + if opt == built_in_defaults.get(section).get(key): del config_sources[section][key] @staticmethod @@ -1371,7 +1362,6 @@ def __getstate__(self): for name in [ "_sections", "is_validated", - "airflow_defaults", ] } @@ -1441,7 +1431,7 @@ def initialize_config() -> AirflowConfigParser: default_config = _parameterized_config_from_template("default_airflow.cfg") - local_conf = AirflowConfigParser(default_config=default_config) + local_conf = AirflowConfigParser() if local_conf.getboolean("core", "unit_test_mode"): # Load test config only diff --git a/airflow/providers/google/config_templates/default_config.cfg b/airflow/providers/google/config_templates/default_config.cfg index 8b9742256465d..bc452d40a97ed 100644 --- a/airflow/providers/google/config_templates/default_config.cfg +++ b/airflow/providers/google/config_templates/default_config.cfg @@ -31,4 +31,4 @@ # Options for google provider # Sets verbose logging for google provider -verbose_logging = False +# verbose_logging = False diff --git a/airflow/utils/db.py b/airflow/utils/db.py index f2b16c5544080..57f6d8182b7cd 100644 --- a/airflow/utils/db.py +++ b/airflow/utils/db.py @@ -33,6 +33,7 @@ import airflow from airflow import settings +from airflow.config_templates.built_in_defaults import built_in_defaults from airflow.configuration import conf from airflow.exceptions import AirflowException from airflow.models import import_all_models @@ -875,8 +876,8 @@ def log_template_exists(): # we will seed the table with the values from pre 2.3.0, so old logs will # still be retrievable. if not session.query(LogTemplate.id).first(): - is_default_log_id = elasticsearch_id == conf.airflow_defaults.get("elasticsearch", "log_id_template") - is_default_filename = filename == conf.airflow_defaults.get("logging", "log_filename_template") + is_default_log_id = elasticsearch_id == built_in_defaults.get("elasticsearch").get("log_id_template") + is_default_filename = filename == built_in_defaults.get("logging").get("log_filename_template") if is_default_log_id and is_default_filename: session.add( LogTemplate( diff --git a/scripts/ci/pre_commit/common_precommit_utils.py b/scripts/ci/pre_commit/common_precommit_utils.py index aef6bc3dcea36..36c8bed34c2a5 100644 --- a/scripts/ci/pre_commit/common_precommit_utils.py +++ b/scripts/ci/pre_commit/common_precommit_utils.py @@ -19,9 +19,14 @@ import hashlib import os import re +from functools import lru_cache from pathlib import Path +from typing import Any + +import jinja2 AIRFLOW_SOURCES_ROOT = Path(__file__).parents[3].resolve() +AIRFLOW_BREEZE_SOURCES_PATH = AIRFLOW_SOURCES_ROOT / "dev" / "breeze" def filter_out_providers_on_non_main_branch(files: list[str]) -> list[str]: @@ -60,3 +65,82 @@ def get_directory_hash(directory: Path, skip_path_regexp: str | None = None) -> if file.is_file() and not file.name.startswith("."): sha.update(file.read_bytes()) return sha.hexdigest() + + +@lru_cache(maxsize=None) +def black_mode(): + from black import Mode, parse_pyproject_toml, target_version_option_callback + + config = parse_pyproject_toml(AIRFLOW_BREEZE_SOURCES_PATH / "pyproject.toml") + + target_versions = set( + target_version_option_callback(None, None, tuple(config.get("target_version", ()))), + ) + + return Mode( + target_versions=target_versions, + line_length=config.get("line_length", Mode.line_length), + is_pyi=bool(config.get("is_pyi", Mode.is_pyi)), + string_normalization=not bool(config.get("skip_string_normalization", not Mode.string_normalization)), + preview=bool(config.get("preview", Mode.preview)), + ) + + +def black_format(content) -> str: + from black import format_str + + return format_str(content, mode=black_mode()) + + +def render_template_from_file( + searchpath: Path, + template_name: str, + context: dict[str, Any], + extension: str, + autoescape: bool = True, + keep_trailing_newline: bool = False, +) -> str: + """ + Renders template based on its name. Reads the template from _TEMPLATE.md.jinja2 in current dir. + :param searchpath: Path to search images in + :param template_name: name of the template to use + :param context: Jinja2 context + :param extension: Target file extension + :param autoescape: Whether to autoescape HTML + :param keep_trailing_newline: Whether to keep the newline in rendered output + :return: rendered template + """ + template_loader = jinja2.FileSystemLoader(searchpath=searchpath) + template_env = jinja2.Environment( + loader=template_loader, + undefined=jinja2.StrictUndefined, + autoescape=autoescape, + keep_trailing_newline=keep_trailing_newline, + ) + template = template_env.get_template(f"{template_name}_TEMPLATE{extension}.jinja2") + return template.render(context) + + +def render_template_from_string( + template_string: str, + context: dict[str, Any] | None = None, + autoescape: bool = True, + keep_trailing_newline: bool = False, +) -> str: + """ + Renders template from string. + :param template_string: template string + :param context: Jinja2 context + :param autoescape: Whether to autoescape HTML + :param keep_trailing_newline: Whether to keep the newline in rendered output + :return: rendered template + """ + + template_env = jinja2.Environment( + undefined=jinja2.StrictUndefined, + autoescape=autoescape, + keep_trailing_newline=keep_trailing_newline, + ).from_string(template_string) + if context is None: + context = {} + return template_env.render(context) diff --git a/scripts/ci/pre_commit/pre_commit_check_pre_commit_hooks.py b/scripts/ci/pre_commit/pre_commit_check_pre_commit_hooks.py index 2cee66a4b28b3..c8279ac330255 100755 --- a/scripts/ci/pre_commit/pre_commit_check_pre_commit_hooks.py +++ b/scripts/ci/pre_commit/pre_commit_check_pre_commit_hooks.py @@ -28,20 +28,23 @@ sys.path.insert(0, str(Path(__file__).parent.resolve())) # make sure common_precommit_utils is imported from collections import defaultdict # noqa: E402 -from functools import lru_cache # noqa: E402 from typing import Any # noqa: E402 import yaml # noqa: E402 -from common_precommit_utils import insert_documentation # noqa: E402 +from common_precommit_utils import ( # noqa: E402 + AIRFLOW_BREEZE_SOURCES_PATH, + AIRFLOW_SOURCES_ROOT, + black_format, + insert_documentation, + render_template_from_file, +) from rich.console import Console # noqa: E402 from tabulate import tabulate # noqa: E402 console = Console(width=400, color_system="standard") -AIRFLOW_SOURCES_PATH = Path(__file__).parents[3].resolve() -AIRFLOW_BREEZE_SOURCES_PATH = AIRFLOW_SOURCES_PATH / "dev" / "breeze" PRE_COMMIT_IDS_PATH = AIRFLOW_BREEZE_SOURCES_PATH / "src" / "airflow_breeze" / "pre_commit_ids.py" -PRE_COMMIT_YAML_FILE = AIRFLOW_SOURCES_PATH / ".pre-commit-config.yaml" +PRE_COMMIT_YAML_FILE = AIRFLOW_SOURCES_ROOT / ".pre-commit-config.yaml" def get_errors_and_hooks(content: Any, max_length: int) -> tuple[list[str], dict[str, list[str]], list[str]]: @@ -75,67 +78,10 @@ def get_errors_and_hooks(content: Any, max_length: int) -> tuple[list[str], dict return errors, hooks, image_hooks -def render_template( - searchpath: Path, - template_name: str, - context: dict[str, Any], - extension: str, - autoescape: bool = True, - keep_trailing_newline: bool = False, -) -> str: - """ - Renders template based on its name. Reads the template from _TEMPLATE.md.jinja2 in current dir. - :param searchpath: Path to search images in - :param template_name: name of the template to use - :param context: Jinja2 context - :param extension: Target file extension - :param autoescape: Whether to autoescape HTML - :param keep_trailing_newline: Whether to keep the newline in rendered output - :return: rendered template - """ - import jinja2 - - template_loader = jinja2.FileSystemLoader(searchpath=searchpath) - template_env = jinja2.Environment( - loader=template_loader, - undefined=jinja2.StrictUndefined, - autoescape=autoescape, - keep_trailing_newline=keep_trailing_newline, - ) - template = template_env.get_template(f"{template_name}_TEMPLATE{extension}.jinja2") - content: str = template.render(context) - return content - - -@lru_cache(maxsize=None) -def black_mode(): - from black import Mode, parse_pyproject_toml, target_version_option_callback - - config = parse_pyproject_toml(AIRFLOW_BREEZE_SOURCES_PATH / "pyproject.toml") - - target_versions = set( - target_version_option_callback(None, None, tuple(config.get("target_version", ()))), - ) - - return Mode( - target_versions=target_versions, - line_length=bool(config.get("line_length", Mode.line_length)), - is_pyi=bool(config.get("is_pyi", Mode.is_pyi)), - string_normalization=not bool(config.get("skip_string_normalization", not Mode.string_normalization)), - preview=bool(config.get("preview", Mode.preview)), - ) - - -def black_format(content) -> str: - from black import format_str - - return format_str(content, mode=black_mode()) - - def prepare_pre_commit_ids_py_file(pre_commit_ids): PRE_COMMIT_IDS_PATH.write_text( black_format( - render_template( + render_template_from_file( searchpath=AIRFLOW_BREEZE_SOURCES_PATH / "src" / "airflow_breeze", template_name="pre_commit_ids", context={"PRE_COMMIT_IDS": pre_commit_ids}, @@ -159,7 +105,7 @@ def update_static_checks_array(hooks: dict[str, list[str]], image_hooks: list[st rows.append((hook_id, formatted_hook_description, " * " if hook_id in image_hooks else " ")) formatted_table = "\n" + tabulate(rows, tablefmt="grid", headers=("ID", "Description", "Image")) + "\n\n" insert_documentation( - file_path=AIRFLOW_SOURCES_PATH / "STATIC_CODE_CHECKS.rst", + file_path=AIRFLOW_SOURCES_ROOT / "STATIC_CODE_CHECKS.rst", content=formatted_table.splitlines(keepends=True), header=" .. BEGIN AUTO-GENERATED STATIC CHECK LIST", footer=" .. END AUTO-GENERATED STATIC CHECK LIST", diff --git a/scripts/ci/pre_commit/pre_commit_update_common_sql_api_stubs.py b/scripts/ci/pre_commit/pre_commit_update_common_sql_api_stubs.py index 5a85b84a0c844..4ab041e2cc808 100755 --- a/scripts/ci/pre_commit/pre_commit_update_common_sql_api_stubs.py +++ b/scripts/ci/pre_commit/pre_commit_update_common_sql_api_stubs.py @@ -25,18 +25,19 @@ import jinja2 import sys import textwrap -from functools import lru_cache from pathlib import Path from rich.console import Console +sys.path.insert(0, str(Path(__file__).parent.resolve())) # make sure common_precommit_utils is imported +from common_precommit_utils import black_format, AIRFLOW_SOURCES_ROOT # noqa: E402 + if __name__ not in ("__main__", "__mp_main__"): raise SystemExit( "This file is intended to be executed as an executable program. You cannot use it as a module." f"To execute this script, run ./{__file__} [FILE] ..." ) -AIRFLOW_SOURCES_ROOT = Path(__file__).parents[3].resolve() PROVIDERS_ROOT = AIRFLOW_SOURCES_ROOT / "airflow" / "providers" COMMON_SQL_ROOT = PROVIDERS_ROOT / "common" / "sql" OUT_DIR = AIRFLOW_SOURCES_ROOT / "out" @@ -46,33 +47,6 @@ console = Console(width=400, color_system="standard") -@lru_cache(maxsize=None) -def black_mode(): - from black import Mode, parse_pyproject_toml, target_version_option_callback - - config = parse_pyproject_toml(os.path.join(AIRFLOW_SOURCES_ROOT, "pyproject.toml")) - - target_versions = set( - target_version_option_callback(None, None, tuple(config.get("target_version", ()))), - ) - - return Mode( - target_versions=target_versions, - line_length=config.get("line_length", Mode.line_length), - is_pyi=bool(config.get("is_pyi", Mode.is_pyi)), - string_normalization=not bool(config.get("skip_string_normalization", not Mode.string_normalization)), - experimental_string_processing=bool( - config.get("experimental_string_processing", Mode.experimental_string_processing) - ), - ) - - -def black_format(content) -> str: - from black import format_str - - return format_str(content, mode=black_mode()) - - class ConsoleDiff(difflib.Differ): def _dump(self, tag, x, lo, hi): """Generate comparison results for a same-tagged range.""" diff --git a/scripts/ci/pre_commit/pre_commit_yaml_to_cfg.py b/scripts/ci/pre_commit/pre_commit_yaml_to_cfg.py index 97cb921a696e9..13f20456c4937 100755 --- a/scripts/ci/pre_commit/pre_commit_yaml_to_cfg.py +++ b/scripts/ci/pre_commit/pre_commit_yaml_to_cfg.py @@ -22,8 +22,23 @@ from __future__ import annotations import os +import sys +from pathlib import Path +from typing import Any import yaml +from rich.console import Console + +sys.path.insert(0, str(Path(__file__).parent.resolve())) # make sure common_precommit_utils is imported +from common_precommit_utils import ( # noqa: E402 + AIRFLOW_SOURCES_ROOT, + black_format, + render_template_from_string, +) + +console = Console(force_terminal=True, color_system="standard", width=400) + +PYTHON_BUILT_IN_DEFAULTS_PATH = AIRFLOW_SOURCES_ROOT / "airflow" / "config_templates" / "built_in_defaults.py" FILE_HEADER = """# # Licensed to the Apache Software Foundation (ASF) under one @@ -55,8 +70,36 @@ # ----------------------- TEMPLATE BEGINS HERE ----------------------- """ +PYTHON_BUILT_IN_DEFAULTS_CONTENT = """# +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you under the Apache License, Version 2.0 (the +# "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, +# software distributed under the License is distributed on an +# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +# KIND, either express or implied. See the License for the +# specific language governing permissions and limitations +# under the License. + +# Those are built-in defaults for Airflow's default configuration. When there +# are no defaults stored in airflow configuration file, those defaults are used +# as a fallback for the configuration values +# +# PLEASE DO NOT MODIFY THIS FILE IT IS AUTOGENERATED FROM THE CONFIGURATION.YAML FILE +from __future__ import annotations + +built_in_defaults: dict[str, dict[str, str]] = {{ config_defaults_dict }} +""" + -def read_default_config_yaml(file_path: str) -> dict: +def read_default_config_yaml(file_path: Path) -> dict: """ Read Airflow configs from YAML file @@ -69,21 +112,53 @@ def read_default_config_yaml(file_path: str) -> dict: return yaml.safe_load(config_file) -def write_config(yaml_config_file_path: str, default_cfg_file_path: str): +def write_config( + yaml_config_file_path: Path, default_cfg_file_path: Path, built_in_defaults_file_path: Path | None +): """ Write config to default_airflow.cfg file - :param yaml_config_file_path: Full path to config.yaml - :param default_cfg_file_path: Full path to default_airflow.cfg + :param yaml_config_file_path: Path to config.yaml + :param default_cfg_file_path: Path to default_airflow.cfg + :param built_in_defaults_file_path: Path to built_in_defaults.py """ - print(f"Converting {yaml_config_file_path} to {default_cfg_file_path}") + console.print(f"[blue]Converting:[/] {yaml_config_file_path} to {default_cfg_file_path}") + config_yaml = read_default_config_yaml(yaml_config_file_path) with open(default_cfg_file_path, "w") as configfile: configfile.writelines(FILE_HEADER) - config_yaml = read_default_config_yaml(yaml_config_file_path) for section_name, section in config_yaml.items(): _write_section(configfile, section_name, section) + if built_in_defaults_file_path: + print(f"[blue]Also preparing:[/] {built_in_defaults_file_path}") + config_defaults_dict = {} + for section in config_yaml: + section_dict: dict[str, Any] = {} + config_defaults_dict[section] = section_dict + options = config_yaml[section]["options"] + for option in options: + if options[option]["default"] is not None: + if ( + option in ["log_filename_template", "log_processor_filename_template"] + and section == "logging" + ): + section_dict[option] = ( + options[option]["default"].replace("{{", "{").replace("}}", "}") + ) + else: + section_dict[option] = options[option]["default"] + built_in_defaults_file_path.write_text( + black_format( + render_template_from_string( + template_string=PYTHON_BUILT_IN_DEFAULTS_CONTENT, + context={"config_defaults_dict": config_defaults_dict}, + autoescape=False, + keep_trailing_newline=True, + ) + ) + ) + def _write_section(configfile, section_name, section): configfile.write(f"\n[{section_name}]\n") @@ -133,20 +208,20 @@ def _write_option(configfile, idx, option_name, option): value = " " + option["default"] else: value = "" - configfile.write(f"{option_name} ={value}\n") + configfile.write(f"# {option_name} ={value}\n") else: configfile.write(f"# {option_name} =\n") if __name__ == "__main__": - airflow_config_dir = os.path.join( - os.path.dirname(__file__), os.pardir, os.pardir, os.pardir, "airflow", "config_templates" - ) - airflow_default_config_path = os.path.join(airflow_config_dir, "default_airflow.cfg") - airflow_config_yaml_file_path = os.path.join(airflow_config_dir, "config.yml") + airflow_config_dir = AIRFLOW_SOURCES_ROOT / "airflow" / "config_templates" + airflow_default_config_path = airflow_config_dir / "default_airflow.cfg" + airflow_config_yaml_file_path = airflow_config_dir / "config.yml" write_config( - yaml_config_file_path=airflow_config_yaml_file_path, default_cfg_file_path=airflow_default_config_path + yaml_config_file_path=airflow_config_yaml_file_path, + default_cfg_file_path=airflow_default_config_path, + built_in_defaults_file_path=Path(airflow_config_dir) / "built_in_defaults.py", ) providers_dir = os.path.join( @@ -160,6 +235,7 @@ def _write_option(configfile, idx, option_name, option): and os.path.isfile(os.path.join(root, "default_config.cfg")) ): write_config( - yaml_config_file_path=os.path.join(root, "config.yml"), - default_cfg_file_path=os.path.join(root, "default_config.cfg"), + yaml_config_file_path=Path(root) / "config.yml", + default_cfg_file_path=Path(root) / "default_config.cfg", + built_in_defaults_file_path=None, ) diff --git a/tests/core/test_configuration.py b/tests/core/test_configuration.py index 6cbe0e9076988..ab2a87c0d84a7 100644 --- a/tests/core/test_configuration.py +++ b/tests/core/test_configuration.py @@ -33,6 +33,7 @@ from pytest import param from airflow import configuration +from airflow.config_templates.built_in_defaults import built_in_defaults from airflow.configuration import ( AirflowConfigException, AirflowConfigParser, @@ -744,13 +745,13 @@ def test_as_dict_respects_sensitive_cmds(self): assert "sql_alchemy_conn" in conf_materialize_cmds["database"] assert "sql_alchemy_conn_cmd" not in conf_materialize_cmds["database"] - if conf_conn == test_conf.airflow_defaults["database"]["sql_alchemy_conn"]: + if conf_conn == built_in_defaults["database"]["sql_alchemy_conn"]: assert conf_materialize_cmds["database"]["sql_alchemy_conn"] == "my-super-secret-conn" assert "sql_alchemy_conn_cmd" in conf_maintain_cmds["database"] assert conf_maintain_cmds["database"]["sql_alchemy_conn_cmd"] == "echo -n my-super-secret-conn" - if conf_conn == test_conf.airflow_defaults["database"]["sql_alchemy_conn"]: + if conf_conn == built_in_defaults["database"]["sql_alchemy_conn"]: assert "sql_alchemy_conn" not in conf_maintain_cmds["database"] else: assert "sql_alchemy_conn" in conf_maintain_cmds["database"]