Fluxos de trabalho gerenciados pela Amazon para Apache Airflow (Amazon MWAA) oferece recursos robustos de orquestração para fluxos de trabalho de dados, mas o gerenciamento de permissões de DAG em escala apresenta desafios operacionais significativos. À medida que as organizações aumentam seus ambientes e equipes de fluxo de trabalho, atribuir e manter manualmente as permissões dos usuários se torna um gargalo que pode afetar a segurança e a produtividade.
As abordagens tradicionais exigem que os administradores configurem manualmente o controle de acesso baseado em função (RBAC) para cada DAG, levando a:
- Atribuições de permissão inconsistentes entre equipes
- Provisionamento de acesso atrasado para novos membros da equipe
- Aumento do risco de erro humano no gerenciamento de permissões
- Sobrecarga operacional significativa que não pode ser dimensionada
Há outra maneira de fazer isso definindo funções RBAC personalizadas conforme mencionado neste Guia do usuário do Amazon MWAA. Porém, não usa Tags de fluxo de ar para fazer isso.
Nesta postagem, mostramos como usar tags Apache Airflow para gerenciar sistematicamente permissões de DAG, reduzindo a carga operacional e mantendo controles de segurança robustos que complementam as medidas de segurança em nível de infraestrutura.
Pré-requisitos
Para implementar esta solução, você precisa de:
Recursos da AWS:
- Um ambiente Amazon MWAA (versão 2.7.2 ou posterior, sem suporte no Airflow 3.0)
- Funções do IAM configuradas para acesso ao Amazon MWAA com relações de confiança apropriadas
- Bucket Amazon Easy Storage Service (Amazon S3) para armazenamento Amazon MWAA DAG com permissões adequadas
Permissões:
- Permissões do IAM para criar e modificar tokens de login na internet do Amazon MWAA
- Função de execução do Amazon MWAA com permissões para acessar o banco de dados de metadados Apache Airflow
- Acesso administrativo para configurar funções e permissões do Apache Airflow
Visão geral da solução
O sistema automatizado de gerenciamento de permissões consiste em quatro componentes principais que trabalham juntos para fornecer controle de acesso seguro e escalável. O diagrama a seguir mostra o fluxo de trabalho de como a solução funciona.

- Camada de integração IAM – As funções do AWS IAM são mapeadas diretamente para as funções do Apache Airflow. Em seguida, os usuários são autenticados por meio do AWS IAM e recebem automaticamente as funções correspondentes do Airflow. Isso oferece suporte a funções de usuário individuais e padrões de acesso baseados em grupo.
Observação:- O controle de acesso baseado em IAM ao Amazon MWAA funciona para funções padrão do Apache Airflow. Para funções personalizadas, o usuário administrador pode atribuir a função personalizada usando a UI do Apache Airflow, conforme mencionado no Postagem do Centro de Conhecimento e no Guia do usuário do Amazon MWAA.
- Se estiver usando outros autenticadores, as permissões do DAG baseadas em tags continuarão a funcionar conforme indicado no Weblog de Huge Knowledge da AWS publicar.
- Configuração baseada em tags – Tags Apache Airflow definidas em DAGs são usadas para declarar requisitos de acesso. Ele oferece suporte a permissões somente leitura, edição e exclusão.
- Mecanismo de sincronização automatizado – O DAG programado verifica todos os DAGs ativos em busca de tags de permissão com base na programação CRON. Em seguida, ele processa tags e atualiza as permissões RBAC do Apache Airflow de acordo. Em seguida, ele fornece uma configuração baseada para controlar a limpeza das permissões existentes.
- Aplicação de controle de acesso baseado em função – O Apache Airflow RBAC impõe as permissões configuradas armazenando nas tabelas de metadados de funções e permissões do Apache Airflow. Os usuários veem apenas os DAGs aos quais têm acesso. Eles têm controle granular sobre a leitura em comparação com as permissões de edição.
Fluxo de dados
- O usuário do Amazon MWAA assume uma função do IAM para acessar a UI do Amazon MWAA.
- O desenvolvedor do DAG adiciona tags relevantes à definição do DAG.
manage_dag_permissionsO DAG implantado no ambiente Amazon MWAA é executado em uma programação CRON, por exemplo, diário.- O DAG atualiza as respectivas permissões de função para o DAG atualizando os metadados do Apache Airflow no banco de dados Apache Airflow.
- Os usuários ganham ou perdem acesso com base nas funções atribuídas.
Nossa solução se baseia na integração IAM existente do Amazon MWAA, ao mesmo tempo em que amplia a funcionalidade por meio de automação personalizada:
- Autenticação e mapeamento de funções – Os usuários são autenticados por meio de funções do AWS IAM que são mapeadas diretamente para funções correspondentes do Airflow.
- Criação automatizada de usuários – No primeiro login, os usuários são criados automaticamente no banco de dados de metadados do Apache Airflow com atribuições de funções apropriadas.
- Controle de permissão baseado em tags – Cada função do Apache Airflow contém permissões específicas de DAG com base em tags definidas nos DAGs.
- Sincronização automatizada – Um script agendado mantém permissões à medida que os DAGs são adicionados ou modificados.
Etapa 1: configurar o mapeamento de funções do IAM para o Airflow
Primeiro, estabeleça o mapeamento entre os principais do IAM e as funções do Apache Airflow. Para conceder permissão usando o AWS Administration Console, conclua as seguintes etapas:
- Faça login em sua conta AWS e abra o console IAM.
- No painel de navegação esquerdo, escolha Usuáriose escolha seu usuário IAM do Amazon MWAA na tabela de usuários.
- No detalhes do usuário página, abaixo Resumoescolha o Permissões guia e escolha Políticas de permissões para expandir o cartão e escolher Adicionar permissões.
- No Conceder permissões seção, escolha Anexe políticas existentes diretamentee escolha Criar política para criar e anexar sua própria política de permissões personalizadas.
- No Criar política página, escolha JSONe copie e cole a seguinte política de permissões JSON no editor de políticas. Esta política concede acesso ao servidor internet ao usuário com a função Public Apache Airflow padrão.
{
"Model": "2012-10-17",
"Assertion": (
{
"Impact": "Enable",
"Motion": "airflow:CreateWebLoginToken",
"Useful resource": "arn:aws:airflow:area:account-id:atmosphere/your-environment-name"
}
)
}Etapa 2: criar o DAG de gerenciamento automatizado de permissões
Agora, crie um DAG que gerenciará automaticamente as permissões com base em tags.
from airflow import DAG, settings
from airflow.operators.python import PythonOperator
from sqlalchemy import textual content
import pendulum
import logging
dag_id = "manage_dag_permissions"
class Constants:
"""
Constants class to carry fixed values used all through the code.
"""
AB_VIEW_MENU = "ab_view_menu"
AB_PERMISSION = "ab_permission"
AB_ROLE = "ab_role"
AB_PERMISSION_VIEW = "ab_permission_view"
AB_PERMISSION_VIEW_ROLE = "ab_permission_view_role"
DAG_TAG = "dag_tag"
CAN_READ = "can_read"
CAN_EDIT = "can_edit"
CAN_DELETE = "can_delete"
def _execute_query(sql_text, params=None, fetch=True):
"""
Execute a parameterized SQL question in opposition to the Airflow metadata DB.
All queries use SQLAlchemy textual content() with bind parameters to forestall SQL injection.
Parameters:
sql_text: SQL string with :named bind parameters
params: dict of parameter values
fetch: If True, return listing of first-column values; if False, commit and return None
Returns:
Checklist of values (first column) if fetch=True, else None
Raises:
Re-raises any exception after rollback and logging
"""
session = settings.Session()
strive:
stmt = textual content(sql_text)
if fetch:
consequence = session.execute(stmt, params or {}).fetchall()
return (row(0) for row in consequence)
else:
session.execute(stmt, params or {})
session.commit()
return None
besides Exception as e:
session.rollback()
logging.error(f"DB question error (fetch={fetch}): {kind(e).__name__}: {e}")
elevate
lastly:
session.shut()
def fetch_airflow_role_id(role_name):
"""
Fetch function id of a given function identify utilizing parameterized question.
"""
consequence = _execute_query(
"SELECT id FROM ab_role WHERE identify = :role_name",
{"role_name": role_name},
)
if not consequence:
elevate ValueError(f"Airflow function not discovered: {role_name}")
logging.information("Fetched function ID efficiently")
return consequence(0)
def fetch_airflow_permission_id(permission_name):
"""
Fetch permission id of a given permission utilizing parameterized question.
"""
consequence = _execute_query(
"SELECT id FROM ab_permission WHERE identify = :perm_name",
{"perm_name": permission_name},
)
if not consequence:
elevate ValueError(f"Airflow permission not discovered: {permission_name}")
logging.information("Fetched permission ID efficiently")
return consequence(0)
def fetch_airflow_menu_object_ids(dag_names):
"""
Fetch view_menu IDs for an inventory of DAG useful resource names.
Makes use of parameterized IN-clause through particular person bind params.
Parameters:
dag_names: listing of DAG useful resource names (e.g. ('DAG:my_dag1', 'DAG:my_dag2'))
Returns:
listing of view_menu IDs
"""
if not dag_names:
return ()
# Construct parameterized IN clause: :p0, :p1, :p2, ...
param_names = (f":p{i}" for i in vary(len(dag_names)))
params = {f"p{i}": identify for i, identify in enumerate(dag_names)}
in_clause = ", ".be part of(param_names)
consequence = _execute_query(
f"SELECT id FROM ab_view_menu WHERE identify IN ({in_clause})",
params,
)
logging.information(f"Fetched {len(consequence)} view menu IDs")
return consequence
def fetch_perms_obj_association_ids(perm_id, view_menu_ids):
"""
Fetch permission_view IDs for a permission and listing of view_menu IDs.
Makes use of parameterized question.
"""
if not view_menu_ids:
return ()
param_names = (f":vm{i}" for i in vary(len(view_menu_ids)))
params = {f"vm{i}": vm_id for i, vm_id in enumerate(view_menu_ids)}
params("perm_id") = perm_id
in_clause = ", ".be part of(param_names)
consequence = _execute_query(
f"SELECT id FROM ab_permission_view WHERE permission_id = :perm_id AND view_menu_id IN ({in_clause})",
params,
)
logging.information(f"Fetched {len(consequence)} permission-view affiliation IDs")
return consequence
def fetch_dag_ids_by_tag(tag_name):
"""
Fetch DAG IDs with a given tag identify utilizing parameterized question.
"""
consequence = _execute_query(
"SELECT DISTINCT dag_id FROM dag_tag WHERE identify = :tag_name",
{"tag_name": tag_name},
)
logging.information(f"Fetched {len(consequence)} DAG IDs for tag")
return consequence
def associate_permission_to_object(perm_id, view_menu_ids):
"""
Affiliate permission to view_menu objects (DAGs) utilizing parameterized INSERT.
"""
session = settings.Session()
strive:
for vm_id in view_menu_ids:
session.execute(
textual content(
"INSERT INTO ab_permission_view (permission_id, view_menu_id) "
"VALUES (:perm_id, :vm_id) "
"ON CONFLICT (permission_id, view_menu_id) DO NOTHING"
),
{"perm_id": perm_id, "vm_id": vm_id},
)
session.commit()
logging.information(f"Related permission to {len(view_menu_ids)} view menus")
besides Exception as e:
session.rollback()
logging.error(f"Error associating permission to things: {kind(e).__name__}: {e}")
elevate
lastly:
session.shut()
def associate_permission_to_role(permission_view_ids, role_id):
"""
Affiliate permission_view entries to a task utilizing parameterized INSERT.
"""
session = settings.Session()
strive:
for pv_id in permission_view_ids:
session.execute(
textual content(
"INSERT INTO ab_permission_view_role (permission_view_id, role_id) "
"VALUES (:pv_id, :role_id) "
"ON CONFLICT (permission_view_id, role_id) DO NOTHING"
),
{"pv_id": pv_id, "role_id": role_id},
)
session.commit()
logging.information(f"Related {len(permission_view_ids)} permissions to function")
besides Exception as e:
session.rollback()
logging.error(f"Error associating permissions to function: {kind(e).__name__}: {e}")
elevate
lastly:
session.shut()
def validate_if_permission_granted(permission_view_ids, role_id):
"""
Validate if given permissions are related to given function utilizing parameterized question.
"""
if not permission_view_ids:
return ()
param_names = (f":pv{i}" for i in vary(len(permission_view_ids)))
params = {f"pv{i}": pv_id for i, pv_id in enumerate(permission_view_ids)}
params("role_id") = role_id
in_clause = ", ".be part of(param_names)
consequence = _execute_query(
f"SELECT id FROM ab_permission_view_role "
f"WHERE permission_view_id IN ({in_clause}) AND role_id = :role_id",
params,
)
logging.information(f"Validated {len(consequence)} permission grants")
return consequence
def clean_up_existing_dag_permissions_for_role(role_id):
"""
Clear up current DAG permissions for a given function utilizing parameterized question.
Notice: this creates a short window the place the function has no DAG permissions.
"""
_execute_query(
"DELETE FROM ab_permission_view_role WHERE id IN ("
" SELECT pvr.id"
" FROM ab_permission_view_role pvr"
" INNER JOIN ab_permission_view pv ON pvr.permission_view_id = pv.id"
" INNER JOIN ab_view_menu vm ON pv.view_menu_id = vm.id"
" WHERE pvr.role_id = :role_id AND vm.identify LIKE :dag_prefix"
")",
{"role_id": role_id, "dag_prefix": "DAG:%"},
fetch=False,
)
logging.information("Cleaned up current DAG permissions for function")
def sync_permission(config_data):
"""
Sync permissions based mostly on the config.
Parameters:
config_data: dict with keys:
- airflow_role_name: identify of the customized Airflow function
- managed_dags: listing of DAG IDs to grant full administration permissions on
(can_read, can_edit, can_delete)
- do_cleanup: if True, take away all current DAG:* permissions first
"""
# Get the function ID for function identify
role_id = fetch_airflow_role_id(config_data("airflow_role_name"))
# Clear up current DAG stage permissions if requested
if config_data.get("do_cleanup", False):
clean_up_existing_dag_permissions_for_role(role_id)
managed_dags = config_data.get("managed_dags", ())
if not managed_dags:
logging.information("No managed DAGs discovered, skipping permission sync")
return
# Decide which permissions to grant (default: can_read solely)
permissions = config_data.get("permissions", (Constants.CAN_READ))
# Construct DAG useful resource names (e.g. ("DAG:my_dag1", "DAG:my_dag2"))
dag_resource_names = (f"DAG:{dag.strip()}" for dag in managed_dags)
# Get IDs for DAG view_menu objects
vm_ids = fetch_airflow_menu_object_ids(dag_resource_names)
if not vm_ids:
logging.information("No view_menu entries discovered for managed DAGs")
return
# Grant the configured permissions on every managed DAG
all_perm_view_ids = ()
for perm_name in permissions:
perm_id = fetch_airflow_permission_id(perm_name)
associate_permission_to_object(perm_id, vm_ids)
all_perm_view_ids += fetch_perms_obj_association_ids(perm_id, vm_ids)
# Affiliate permission_view entries with the function and validate
if all_perm_view_ids and role_id:
associate_permission_to_role(all_perm_view_ids, role_id)
validate_if_permission_granted(all_perm_view_ids, role_id)
def sync_permissions_with_tags(role_mappings):
"""
For every function mapping, fetch DAG IDs by tag and sync permissions.
"""
for role_map in role_mappings:
username = listing(role_map.keys())(0)
airflow_role = role_map(username)("airflow_role")
edit_tag_name = role_map(username)("airflow_edit_tag")
config_data = {
"airflow_role_name": airflow_role,
"managed_dags": fetch_dag_ids_by_tag(edit_tag_name),
"permissions": role_map(username).get("permissions", (Constants.CAN_READ)),
"do_cleanup": role_map(username).get("do_cleanup", True),
}
logging.information(f"Syncing permissions for airflow function")
sync_permission(config_data)
logging.information("Accomplished permission sync for function")
"""
Add new roles and permissions right here.
Format:
{
"": {
"airflow_role": ,
"airflow_edit_tag": ,
"permissions": ,
"do_cleanup":
}
},
IMPORTANT - ROLE SETUP:
When creating a brand new customized function (e.g. "analytics_reporting", "marketing_analyst")
within the Airflow UI (Safety > Checklist Roles), you MUST copy the Viewer function's
permissions into the brand new function. The Viewer permissions present base UI entry
(browse DAGs, view logs, menu entry, and so on.). --or-- Assign the viewer function as properly.
With out them, customers assigned to
the customized function will be unable to log in to the Airflow UI.
This DAG manages DAG-level permissions on DAG:xxx assets.
Which permissions are granted is managed by the "permissions" listing
in every config entry (choices: can_read, can_edit, can_delete).
It does NOT handle base UI permissions — these should be arrange manually
when creating the function.
Steps to create a brand new customized function:
1. Go to Safety > Checklist Roles > + (Add)
2. Identify it to match the "airflow_role" worth within the config under
3. Copy all permissions from the "Viewer" function into the brand new function
4. Save — this DAG will then robotically add DAG-specific permissions
(as configured within the "permissions" listing) for every tagged DAG
"""
role_mappings = (
{
"analytics_reporting": {
"airflow_role": "analytics_reporting",
"airflow_edit_tag": "analytics_reporting_edit",
"permissions": ("can_read", "can_edit", "can_delete"),
"do_cleanup": True,
}
},
{
"marketing_analyst": {
"airflow_role": "marketing_analyst",
"airflow_edit_tag": "marketing_analyst_edit",
"permissions": ("can_read", "can_edit", "can_delete"),
"do_cleanup": True,
},
},
)
with DAG(
dag_id=dag_id,
schedule="*/15 * * * *",
catchup=False,
start_date=pendulum.datetime(2026, 1, 1, tz="UTC"),
) as dag:
sync_dag_permissions_task = PythonOperator(
task_id="sync_dag_permissions",
python_callable=sync_permissions_with_tags,
op_kwargs={"role_mappings": role_mappings},
)
Etapa 3: marque seus DAGs para controle de acesso
Adicione tags apropriadas aos seus DAGs para especificar quais funções devem ter acesso. As tags são usadas para definir quais funções têm acesso aos DAGs marcados.
# Instance DAG for analytics_reporting
with DAG(
"analytics_reporting_dag",
description="Every day analytics reporting pipeline",
schedule_interval="@each day",
start_date=pendulum.datetime(2023, 1, 1, tz="UTC"),
catchup=False,
tags=("reporting", "analytics", "analytics_reporting_edit")
) as dag:
# DAG duties right here
go
# Instance DAG for marketing_analyst
with DAG(
"marketing_analyst_dag",
description="Every day advertising lead evaluation pipeline",
schedule_interval="@each day",
start_date=pendulum.datetime(2023, 1, 1, tz="UTC"),
catchup=False,
tags=("advertising", "analytics", "marketing_analyst_edit")
) as dag:
# DAG duties right here
goNeste exemplo:
- O
analytics_reportinga função personalizada terá acesso de leitura, edição e exclusão ao DAGanalytics_reporting_dag(e outros DAGs marcados comanalytics_reporting_edit) - O
marketing_analysta função personalizada terá acesso de leitura, edição e exclusão ao DAGmarketing_analyst_dag(e outros DAGs marcados commarketing_analyst_edit)
As permissões exatas concedidas (can_read, can_edit, can_delete) são configuráveis por função no role_mappings config dentro do DAG de gerenciamento de permissões:
role_mappings = (
{
"analytics_reporting": {
"airflow_role": "analytics_reporting",
"airflow_edit_tag": "analytics_reporting_edit",
"permissions": ("can_read", "can_edit", "can_delete"),
"do_cleanup": True,
}
},
{
"marketing_analyst": {
"airflow_role": "marketing_analyst",
"airflow_edit_tag": "marketing_analyst_edit",
"permissions": ("can_read", "can_edit", "can_delete"),
"do_cleanup": True,
},
},
)Observação: antes que esse DAG possa gerenciar permissões para uma função personalizada, a função deve ser criada manualmente na UI do Apache Airflow (Segurança > Listar funções) com as permissões da função de visualizador copiadas. Consulte a Etapa 2 para obter detalhes.
Etapa 4: implantar e testar
- Faça add do DAG de gerenciamento de permissões e dos DAGs marcados para o bucket S3 do seu ambiente Amazon MWAA.
- Aguarde até que o Amazon MWAA detecte e processe os novos DAGs.
- Verifique se o DAG de gerenciamento de permissões é executado com êxito.
- Teste o acesso com diferentes funções de usuário para confirmar a aplicação de permissão adequada.
- Os usuários também podem integre isso aos seus processos de CI/CD.
Solução de problemas
Nesta seção, abordamos alguns problemas comuns e como solucioná-los.
Falhas de sincronização de permissão
Sintoma: o DAG de sincronização de permissão falha com erros de banco de dados.
Causa: permissões insuficientes na função de execução MWAA.
Solução: certifique-se de que a função de execução tenha airflow:CreateWebLoginToken permissão e acesso ao banco de dados.
Tags não sendo processadas
Sintoma: tags DAG estão presentes, mas as permissões não são atualizadas.
Solução: verifique se o DAG está ativo e analisado com sucesso – revise os logs do DAG de sincronização de permissão para erros de processamento.
Os usuários não podem acessar os DAGs esperados
Sintoma: os usuários com funções IAM corretas não podem ver os DAGs
Solução: confirme se o mapeamento de funções do IAM para o Apache Airflow está correto. Verifique se o DAG de sincronização de permissão foi executado com êxito. Verifique o Amazon CloudWatch Logs em busca de erros de atribuição de permissão.
Problemas de desempenho
Sintoma: a sincronização da permissão demora muito ou expira.
Solução: Reduza a frequência de sincronização para ambientes grandes. Considere atualizações de permissão em lote. Monitore o tempo de execução do DAG e otimize adequadamente.
Etapas de depuração
- Verifique a integridade e a conectividade do ambiente Amazon MWAA
- Revise os registros de execução do DAG de sincronização de permissão
- Verifique as configurações de função do IAM e as relações de confiança
- Teste com um único DAG para isolar problemas
- Monitore o CloudWatch Logs para obter mensagens de erro detalhadas
Benefícios e considerações
O gerenciamento automatizado de permissões oferece vantagens operacionais significativas e, ao mesmo tempo, aprimora sua segurança. Você se beneficiará da redução da sobrecarga administrativa à medida que as atribuições manuais de permissões forem removidas, para que você possa escalar perfeitamente, sem encargos adicionais. Sua segurança melhora por meio da aplicação consistente de princípios de privilégio mínimo e da redução de erros humanos. Você aprimorará sua experiência de desenvolvedor com provisionamento de acesso automático que reduz o tempo de integração, enquanto seu sistema oferece suporte a ambientes com mais de 500 DAGs sem degradação de desempenho.
Ao implementar esses sistemas, você deve aderir às principais práticas de segurança. Você deve aplicar o princípio do privilégio mínimo, validar tags para garantir que está processando apenas tags autorizadas e estabelecer mecanismos de auditoria abrangentes, incluindo o registro em log do CloudTrail. Suas medidas de controle de acesso devem restringir as funções de gerenciamento de permissões aos administradores enquanto você utiliza a separação de funções apropriada para diferentes perfis de usuários.
Você precisará considerar diversas limitações técnicas durante sua implementação. O controle de acesso baseado em IAM ao Amazon MWAA funciona apenas com funções padrão do Apache Airflow, e não com funções personalizadas, embora suas permissões baseadas em tags funcionem com autenticadores alternativos. As alterações de permissão são propagadas com base nas programações do DAG, podendo causar atrasos. Você deve estabelecer processos de aprovação para suas alterações de produção, manter o controle de versão para permissões e documentar seus procedimentos de reversão para garantir a resiliência e a segurança do seu sistema.
Limpar
Limpe os recursos após sua experimentação:
- Exclua os ambientes do Amazon MWAA usando o console ou AWS CLI.
- Atualize a política de função do IAM ou exclua a função do IAM, se não for necessário.
Conclusão
Nesta postagem, você aprendeu como automatizar o gerenciamento de permissões DAG no Amazon MWAA usando o sistema de marcação do Apache Airflow. Você viu como implementar o controle de acesso baseado em tags que é dimensionado com eficiência, reduz erros manuais e mantém princípios de segurança com privilégios mínimos em centenas de DAGs. Você também explorou as principais práticas de segurança e considerações técnicas que precisa ter em mente durante a implementação.
Experimente esta solução em seu ambiente Amazon MWAA para agilizar o gerenciamento de permissões. Comece implementando o sistema de marcação em um ambiente de desenvolvimento e, em seguida, implemente-o gradualmente na produção à medida que sua equipe se sentir confortável com a abordagem.