news 2026/9/13 11:48:12

Apache Airflow 中 Cassandra 连接的配置指南:从 Connection 字段到 Hook 源码级解析

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Apache Airflow 中 Cassandra 连接的配置指南:从 Connection 字段到 Hook 源码级解析

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(",")拆分成列表再传给Clustercontact_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-driverCluster级参数。根据 cassandra.rst 文档,除标准 Python 参数外,Hook 额外支持以下参数:

Extra 参数说明
load_balancing_policy指定负载均衡策略,可选RoundRobinPolicy(默认)、DCAwareRoundRobinPolicyAllowListRoundRobinPolicyTokenAwarePolicy
load_balancing_policy_args上述负载均衡策略对应的构造参数
cql_version指定 Cassandra 的 CQL 版本
protocol_version指定使用的原生协议最大版本号
ssl_options当 Cassandra 启用了 SSL 时,指定 SSL 相关细节

在 CassandraHook.init中,cql_versionssl_optionsprotocol_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:子策略名,只能是RoundRobinPolicyDCAwareRoundRobinPolicyWhiteListRoundRobinPolicy三者之一,默认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_policyload_balancing_policy_args

默认负载均衡策略为RoundRobinPolicy。以下是三种常用策略的完整配置:

DCAwareRoundRobinPolicylocal_dcused_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" } }

AllowListRoundRobinPolicyhosts为必填数组):

{ "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-2Port=9042Schema=test_keyspace。测试断言Cluster.__init__收到的参数为contact_points=["host-1", "host-2"]port=9042protocol_version=4load_balancing_policyTokenAwarePolicy实例,可作为你配置时的对照基准。

源码视角:Hook 如何组装并管理连接

连接组装流程

CassandraHook.__init__的执行流程可归纳为:

  1. 通过self.get_connection(cassandra_conn_id)读取连接(默认cassandra_default);
  2. conn.host按逗号拆分为contact_points
  3. conn.port转为int作为port
  4. 若存在login,用PlainTextAuthProvider(login, password)设置auth_provider
  5. 解析load_balancing_policy/load_balancing_policy_args生成负载均衡策略;
  6. 透传cql_versionssl_optionsprotocol_version
  7. Cluster(**conn_config)创建集群对象,并将conn.schema保存为 keyspace;
  8. 会话(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 后端等方式注入连接配置。

注意事项与常见坑

  1. 策略名差异:AllowListRoundRobinPolicyvsWhiteListRoundRobinPolicy:文档中列出的策略名为AllowListRoundRobinPolicy,而当前仓库 CassandraHook 实际匹配并实例化的是 driver 的WhiteListRoundRobinPolicy。配置时请以你安装的cassandra-driver版本实际支持的策略名为准,load_balancing_policy传入未识别的名称不会报错,而是静默回退为默认RoundRobinPolicy——这可能掩盖配置错误,建议配置后通过日志确认实际生效的策略。

  2. WhiteListRoundRobinPolicy缺少hosts会抛错:这是四种策略中唯一会在缺少参数时抛出ValueError的策略(cassandra.py),且该错误在TokenAwarePolicy作为子策略时会由get_lb_policy一并向上抛出。

  3. 端口与端口协议conn.port会被int()转换,若端口字段为空则不会传给Cluster,此时使用cassandra-driver的默认端口 9042。

  4. 安全边界record_exists对表名/keyspace 做了白名单字符校验并对键值使用参数化查询;但在table_exists中表名仅用于元数据查询。请勿将未经验证的用户输入直接拼入 CQL。

  5. 测试覆盖:以上行为均有单元测试(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),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/13 11:43:46

Symfony框架核心特性与PHP企业级开发实践

1. Symfony框架概述与核心特性Symfony是一个基于PHP语言的成熟Web应用框架&#xff0c;自2005年发布以来已成为企业级开发的标准选择。这个全栈框架采用模块化组件设计&#xff0c;其核心思想是"约定优于配置"&#xff0c;同时保持高度的灵活性。与其他PHP框架相比&a…

作者头像 李华
网站建设 2026/9/13 11:43:42

PIC12F675 ADC寄存器配置与采样稳定性实战

简介&#xff1a;本资源是面向嵌入式初学者与PIC单片机开发者的PIC12F675微控制器ADC功能实践例程&#xff0c;聚焦模拟信号采集核心需求&#xff0c;适用于温度检测、电位器读取、传感器数据采集等典型应用场景。压缩包共21个文件&#xff0c;涵盖C源码&#xff08;ADC.C&…

作者头像 李华