Apache Airflow 中 Cassandra 连接的配置指南:从 Connection 字段到 Hook 源码级解析
【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow
Apache Airflow 通过apache-airflow-providers-apache-cassandraProvider 提供对 Apache Cassandra 数据库的连接能力,其核心文档位于 connections/cassandra.rst。本文以该文档为骨架,完整讲解 Cassandra Connection 的字段含义、Extra 高级参数(负载均衡策略、SSL、CQL/协议版本)的配置方法,并结合仓库中 CassandraHook 的源码实现与单元/集成测试,说明这些配置在底层是如何被解析和生效的。读完本文,你将能够独立配置一个生产可用的 Cassandra 连接,并理解为何某些参数会触发策略回退或直接抛错。
Cassandra Connection 概述与默认 Connection ID
Cassandra Connection 是 Airflow 中connection-type: cassandra的连接类型,用于建立与 Apache Cassandra 集群的会话。在 provider.yaml 中可以看到该连接类型与 Hook 的绑定关系:
connection-types: - hook-class-name: airflow.providers.apache.cassandra.hooks.cassandra.CassandraHook hook-name: "Cassandra" connection-type: cassandra默认 Connection ID 为cassandra_default。这一默认值来自 CassandraHook 的类属性:
conn_name_attr = "cassandra_conn_id" default_conn_name = "cassandra_default" conn_type = "cassandra" hook_name = "Cassandra"这意味着:Cassandra Hook 和基于它实现的 Cassandra Operators / Sensors 在不显式指定cassandra_conn_id时,都会自动使用cassandra_default。例如 CassandraTableSensor 和 CassandraRecordSensor 的构造函数中,cassandra_conn_id参数默认值都是CassandraHook.default_conn_name。
配置 Connection 的标准字段
在 Airflow UI(Admin → Connections)或环境变量 / Secret 后端中创建cassandra类型的连接时,需要配置以下字段:
| 字段 | 是否必填 | 说明 |
|---|---|---|
| Host | 必填 | 要连接的 Cassandra 节点主机地址,支持以逗号分隔的多个主机列表,对应 driver 的contact_points |
| Schema | 必填 | 数据库中的 schema(即 keyspace)名称,Hook 会用它建立会话 |
| Login | 必填 | 连接使用的用户名,映射为PlainTextAuthProvider的用户名 |
| Password | 必填 | 连接使用的密码 |
| Port | 必填 | 连接端口(Cassandra 原生协议默认 9042) |
| Extra | 可选 | 以 JSON 字典形式传入的额外参数,用于定制驱动行为 |
值得强调的是Host 支持多主机逗号分隔。从 CassandraHook.init可以看到:
if conn.host: conn_config["contact_points"] = conn.host.split(",")conn.host会先被split(",")拆分成列表再传给Cluster的contact_points。单元测试 test_cassandra.py 验证了这一点:当连接中 host 配置为"127.0.0.1,127.0.0.2"时,cluster.contact_points == ["127.0.0.1", "127.0.0.2"];集成测试 test_cassandra.py 中同样使用host="host-1,host-2"验证Cluster.__init__收到的contact_points=["host-1", "host-2"]。因此,对于多节点集群,直接在 Host 字段中用逗号列出所有种子节点即可,无需写脚本拼接。
Login 与 Password 会一起构造成PlainTextAuthProvider:
if conn.login: conn_config["auth_provider"] = PlainTextAuthProvider(username=conn.login, password=conn.password)Schema 字段则在连接建立后用于cluster.connect(keyspace),即 get_conn() 的实现。
Extra 字段支持的驱动级参数
Extra 字段以 JSON 字典形式承载 Python 官方驱动cassandra-driver的Cluster级参数。根据 cassandra.rst 文档,除标准 Python 参数外,Hook 额外支持以下参数:
| Extra 参数 | 说明 |
|---|---|
load_balancing_policy | 指定负载均衡策略,可选RoundRobinPolicy(默认)、DCAwareRoundRobinPolicy、AllowListRoundRobinPolicy、TokenAwarePolicy |
load_balancing_policy_args | 上述负载均衡策略对应的构造参数 |
cql_version | 指定 Cassandra 的 CQL 版本 |
protocol_version | 指定使用的原生协议最大版本号 |
ssl_options | 当 Cassandra 启用了 SSL 时,指定 SSL 相关细节 |
在 CassandraHook.init中,cql_version、ssl_options、protocol_version三个参数是直接透传给Cluster(**conn_config)的:
cql_version = conn.extra_dejson.get("cql_version", None) if cql_version: conn_config["cql_version"] = cql_version ssl_options = conn.extra_dejson.get("ssl_options", None) if ssl_options: conn_config["ssl_options"] = ssl_options protocol_version = conn.extra_dejson.get("protocol_version", None) if protocol_version: conn_config["protocol_version"] = protocol_version单元测试 test_cassandra.py 展示了同时配置四个参数的完整 Extra 形态,并断言cluster.cql_version == "3.4.4"、cluster.ssl_options == {"ca_certs": "/path/to/certs"}、cluster.protocol_version == 4,证明这些参数确实原样进入了Cluster构造。
负载均衡策略(Load Balancing Policy)详解
四种可用策略
文档列出了四种负载均衡策略,Hook 通过get_lb_policy()静态方法(cassandra.py)按名称实例化:
RoundRobinPolicy(默认):在所有节点间轮询,不感知数据中心拓扑。当 Extra 中未指定策略名或策略名无法识别时,Hook 都会回退到该策略。DCAwareRoundRobinPolicy:优先在本地数据中心内轮询,可配置local_dc(本地数据中心名,默认空串)与used_hosts_per_remote_dc(每个远程数据中心最多使用的节点数,默认0表示全部使用)。AllowListRoundRobinPolicy:只允许访问显式列出的主机。注意:文档中的名称是AllowListRoundRobinPolicy,但当前仓库源码实现使用的是WhiteListRoundRobinPolicy(详见下文注意事项小节)。配置该策略时必须提供hosts列表,否则会抛出ValueError("Hosts must be specified for WhiteListRoundRobinPolicy")。TokenAwarePolicy:基于分区键的 token 感知策略,将请求路由到持有对应数据副本的节点,从而降低网络往返。它需要(可选)指定一个子负载均衡策略:child_load_balancing_policy:子策略名,只能是RoundRobinPolicy、DCAwareRoundRobinPolicy、WhiteListRoundRobinPolicy三者之一,默认RoundRobinPolicy;child_load_balancing_policy_args:子策略的构造参数。
关键实现细节:若TokenAwarePolicy指定的子策略名不在白名单内,Hook 会静默回退到RoundRobinPolicy作为子策略(而非抛错);而未识别的主策略名则回退为顶层RoundRobinPolicy。这一行为被集成测试 test_cassandra.py 明确覆盖:"DoesNotExistPolicy"与非法子策略名均得到RoundRobinPolicy。
策略构造逻辑源码
if policy_name == "DCAwareRoundRobinPolicy": local_dc = policy_args.get("local_dc", "") used_hosts_per_remote_dc = int(policy_args.get("used_hosts_per_remote_dc", 0)) return DCAwareRoundRobinPolicy(local_dc, used_hosts_per_remote_dc) if policy_name == "WhiteListRoundRobinPolicy": hosts = policy_args.get("hosts") if not hosts: raise ValueError("Hosts must be specified for WhiteListRoundRobinPolicy") return WhiteListRoundRobinPolicy(hosts) if policy_name == "TokenAwarePolicy": # child policy 白名单检查,非法时回退 RoundRobinPolicy ...注意used_hosts_per_remote_dc会被int()强制转换为整数,因此即便在 JSON 中写成字符串"SOME_INT_VALUE"也能正确解析——这一细节在 Hook 的 docstring 和集成测试({"used_hosts_per_remote_dc": "3"},见 test_cassandra.py)中均有体现。
Extra 配置示例
示例一:指定ssl_options
当 Cassandra 启用了 SSL 时,在 Extra 中传入一个字典,作为ssl.wrap_socket()的 kwargs。最典型的用法是指定 CA 证书路径:
{ "ssl_options": { "ca_certs": "PATH/TO/CA_CERTS" } }示例二:指定load_balancing_policy与load_balancing_policy_args
默认负载均衡策略为RoundRobinPolicy。以下是三种常用策略的完整配置:
DCAwareRoundRobinPolicy(local_dc与used_hosts_per_remote_dc均可选):
{ "load_balancing_policy": "DCAwareRoundRobinPolicy", "load_balancing_policy_args": { "local_dc": "LOCAL_DC_NAME", "used_hosts_per_remote_dc": "SOME_INT_VALUE" } }AllowListRoundRobinPolicy(hosts为必填数组):
{ "load_balancing_policy": "AllowListRoundRobinPolicy", "load_balancing_policy_args": { "hosts": ["HOST1", "HOST2", "HOST3"] } }TokenAwarePolicy(子策略可选,未指定时默认RoundRobinPolicy):
{ "load_balancing_policy": "TokenAwarePolicy", "load_balancing_policy_args": { "child_load_balancing_policy": "CHILD_POLICY_NAME", "child_load_balancing_policy_args": {} } }组合示例:多主机 + 协议版本 + 负载均衡
参照集成测试 test_cassandra.py,一个生产形态的连接可以这样组织:
{ "load_balancing_policy": "TokenAwarePolicy", "protocol_version": 4 }对应的连接字段为:Host=host-1,host-2、Port=9042、Schema=test_keyspace。测试断言Cluster.__init__收到的参数为contact_points=["host-1", "host-2"]、port=9042、protocol_version=4且load_balancing_policy为TokenAwarePolicy实例,可作为你配置时的对照基准。
源码视角:Hook 如何组装并管理连接
连接组装流程
CassandraHook.__init__的执行流程可归纳为:
- 通过
self.get_connection(cassandra_conn_id)读取连接(默认cassandra_default); - 将
conn.host按逗号拆分为contact_points; - 将
conn.port转为int作为port; - 若存在
login,用PlainTextAuthProvider(login, password)设置auth_provider; - 解析
load_balancing_policy/load_balancing_policy_args生成负载均衡策略; - 透传
cql_version、ssl_options、protocol_version; - 以
Cluster(**conn_config)创建集群对象,并将conn.schema保存为 keyspace; - 会话(Session)惰性创建,首次调用
get_conn()时才执行cluster.connect(self.keyspace)。
会话建立后会被缓存复用(只要session.is_shutdown为假),避免每次调用重复建连。Hook 还提供get_cluster()与shutdown_cluster()方法用于获取集群对象和显式关闭所有连接(cassandra.py)。
两个开箱即用的探测方法
Hook 内置了两个数据探测方法,它们正是两个 Cassandra Sensor 的底层实现:
table_exists(table)(cassandra.py):支持keyspace.table点号写法,通过cluster.metadata检查 keyspace 及其 tables 元数据判断表是否存在;record_exists(table, keys)(cassandra.py):根据keys字典生成SELECT * FROM keyspace.table WHERE k1=%(k1)s AND k2=%(k2)s参数化查询,返回result.one() is not None判断记录是否存在。
值得注意的安全设计:record_exists中的 keyspace 与表名会经过_sanitize_input(cassandra.py)正则校验,只允许\w+字符,否则抛出ValueError。集成测试 test_cassandra.py 的test_possible_sql_injection用"t; DROP TABLE t; SELECT * FROM t"验证了注入防护:这类输入会直接触发ValueError: Invalid input: ...而非执行恶意语句。同时键值通过驱动参数化查询绑定,从根本上避免了 CQL 注入。
与 Sensor 配合使用:等待表与记录就绪
配置好 Connection 后,即可在 DAG 中使用 Cassandra 相关传感器(其使用指南见 operators.rst,前提是必须先配置 Cassandra Connection):
CassandraTableSensor:轮询检查 Cassandra 集群中是否存在指定表,table参数支持点号写法指定 keyspace;CassandraRecordSensor:轮询检查指定表中是否存在满足键值对条件的记录,keys参数以字典形式给出主键-值映射。
系统示例 DAG example_cassandra_dag.py 展示了完整用法:
with DAG( dag_id="example_cassandra_operator", schedule=None, start_date=datetime(2021, 1, 1), default_args={"table": "keyspace_name.table_name"}, catchup=False, tags=["example"], ) as dag: # Replace <table_name> with your actual table name table_sensor = CassandraTableSensor(task_id="cassandra_table_sensor", table="<table_name>") record_sensor = CassandraRecordSensor( task_id="cassandra_record_sensor", keys={"p1": "v1", "p2": "v2"}, table="<table_name>" )要点说明:
table参数已列入template_fields,支持 Jinja 模板渲染(table.py、record.py);- 点号写法用于指定特定 keyspace,例如
table="k.t";若不带 keyspace,则使用 Connection 的 Schema 字段所指定的 keyspace; - 两个 Sensor 的
poke()分别调用hook.table_exists(...)与hook.record_exists(...),并复用cassandra_default连接(可通过cassandra_conn_id覆盖); - 由于两个 Sensor 无依赖关系,可并行执行;实际业务中常将其串接在数据写入任务之后,用于等待下游数据可见。
安装与依赖要求
该 Provider 以独立包分发,安装命令为:
pip install apache-airflow-providers-apache-cassandra根据 README.rst 与 pyproject.toml,当前版本(3.9.5)的要求如下:
- 底层 Airflow 版本
>=2.11.0; - 依赖
apache-airflow-providers-common-compat >=1.8.0; - 依赖
cassandra-driver(Python 3.12 及以上要求>=3.29.2,Python 3.13 要求>=3.29.3,低于 3.12 要求>=3.29.1); - 支持 Python 3.10 ~ 3.13,明确排除了 Python 3.14(见 pyproject.toml 的
requires-python = ">=3.10,!=3.14.*"); - Provider 生命周期状态为
production/ready(provider.yaml)。
安装后,你可以在 Airflow UI 的 Admin → Connections 中新建cassandra类型连接,或通过airflow connections addCLI、环境变量、Secrets 后端等方式注入连接配置。
注意事项与常见坑
策略名差异:
AllowListRoundRobinPolicyvsWhiteListRoundRobinPolicy:文档中列出的策略名为AllowListRoundRobinPolicy,而当前仓库 CassandraHook 实际匹配并实例化的是 driver 的WhiteListRoundRobinPolicy。配置时请以你安装的cassandra-driver版本实际支持的策略名为准,load_balancing_policy传入未识别的名称不会报错,而是静默回退为默认RoundRobinPolicy——这可能掩盖配置错误,建议配置后通过日志确认实际生效的策略。WhiteListRoundRobinPolicy缺少hosts会抛错:这是四种策略中唯一会在缺少参数时抛出ValueError的策略(cassandra.py),且该错误在TokenAwarePolicy作为子策略时会由get_lb_policy一并向上抛出。端口与端口协议:
conn.port会被int()转换,若端口字段为空则不会传给Cluster,此时使用cassandra-driver的默认端口 9042。安全边界:
record_exists对表名/keyspace 做了白名单字符校验并对键值使用参数化查询;但在table_exists中表名仅用于元数据查询。请勿将未经验证的用户输入直接拼入 CQL。测试覆盖:以上行为均有单元测试(unit 测试)与集成测试(integration 测试,需要真实 Cassandra 环境并以
cassandra标记运行)双重保障,修改或排查问题时可直接参考这些用例。
综上,Cassandra Connection 的配置核心是:必填字段保证连接可达与认证,Extra 字段则完全控制驱动的连接行为。理解 CassandraHook 对每个字段的解析路径,就能在配置异常时快速定位是策略回退、参数透传失败还是认证配置问题。
【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考