Apache Airflow 集成 Apache Cassandra:TableSensor 与 RecordSensor 实战指南与源码解析
【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow
Apache Airflow 的 Apache Cassandra Provider 提供了一组传感器(Sensor),用于在工作流中等待 Cassandra 集群中的表(Table)或记录(Record)就绪,从而让 DAG 与 Cassandra 数据写入流程精准对齐。本文以 Apache Cassandra Operators 官方指南 为主线,完整讲解CassandraTableSensor与CassandraRecordSensor的配置方法、参数语义、完整示例 DAG,并结合 CassandraHook 源码 与单元/集成测试,剖析其底层探测原理、连接构建过程与安全细节。读完本文,你将能够在一个真实 Airflow 环境中配置 Cassandra 连接、编写等待表与记录就绪的传感器任务,并理解其内部工作机制。
为什么需要 Cassandra 传感器
Apache Cassandra 是一个开源的分布式 NoSQL 数据库,专为需要线性扩展与高可用性而不牺牲性能的场景而设计。它在商用硬件或云基础设施上提供线性扩展与容错能力,支持多数据中心复制且延迟更低,因此常被用于承载关键任务数据。在数据管道实践中,下游任务常常依赖上游外部系统向 Cassandra 写入特定表或记录:例如,一个 DAG 需要等实时流处理程序把数据刷进某个表后,再触发 ETL 聚合。此时若直接执行查询,往往会因表尚未创建、记录尚未落盘而失败。
Airflow 传感器(Sensor)正是为这种“等待外部条件就绪”的场景设计的:它会在调度周期内反复执行poke()探测,直到条件满足或超时。Apache Cassandra Provider 提供了两个开箱即用的传感器:
| 传感器 | 探测目标 | 关键参数 |
|---|---|---|
CassandraTableSensor | 表中是否存在 | table |
CassandraRecordSensor | 表中是否存在指定记录 | table、keys |
两者都通过CassandraHook与集群交互,并使用点号(dot notation)定位 keyspace 下的表。
前置条件:配置 Cassandra Connection
使用这两个传感器之前,必须先配置一个 Cassandra Connection,完整配置说明见 Apache Cassandra Connection 指南。Cassandra Hook 与 Cassandra 传感器默认使用连接 IDcassandra_default(源码中定义于 cassandra.py 的default_conn_name),你也可以通过cassandra_conn_id参数指定其他连接。
在 Airflow UI 的 Admin → Connections 页面(或通过环境变量、Secrets Backend)创建 Cassandra 连接时,需要填写的字段如下:
- Host(必填):要连接的 Cassandra 主机,支持用逗号分隔的多个主机列表(对应驱动层的
contact_points)。 - Schema(必填):数据库中的 schema,即 keyspace 名称。
- Login(必填):连接用户名。
- Password(必填):连接密码。
- Port(必填):连接端口(默认 9042)。
- Extra(可选):以 JSON 字典形式提供的扩展参数,支持以下标准参数之外的选项:
load_balancing_policy:负载均衡策略,可选RoundRobinPolicy、DCAwareRoundRobinPolicy、AllowListRoundRobinPolicy、TokenAwarePolicy,默认RoundRobinPolicy;load_balancing_policy_args:上述策略的参数;cql_version:指定 Cassandra 的 CQL 版本;protocol_version:指定要使用的原生协议最大版本;ssl_options:当 Cassandra 启用 SSL 时,指定 SSL 相关细节。
Extra 字段配置示例
1. 配置ssl_options:如果 Cassandra 启用了 SSL,可在 Extra 字段中传入ssl.wrap_socket()的 kwargs 字典,例如:
{ "ssl_options": { "ca_certs": "PATH/TO/CA_CERTS" } }2. 配置load_balancing_policy与load_balancing_policy_args:默认负载均衡策略为RoundRobinPolicy,以下是几种常见策略的示例:
DCAwareRoundRobinPolicy(多数据中心感知轮询):
{ "load_balancing_policy": "DCAwareRoundRobinPolicy", "load_balancing_policy_args": { "local_dc": "LOCAL_DC_NAME", "used_hosts_per_remote_dc": "SOME_INT_VALUE" } }AllowListRoundRobinPolicy(白名单轮询,源码实现对应WhiteListRoundRobinPolicy):
{ "load_balancing_policy": "AllowListRoundRobinPolicy", "load_balancing_policy_args": { "hosts": ["HOST1", "HOST2", "HOST3"] } }TokenAwarePolicy(令牌感知,可指定子策略):
{ "load_balancing_policy": "TokenAwarePolicy", "load_balancing_policy_args": { "child_load_balancing_policy": "CHILD_POLICY_NAME", "child_load_balancing_policy_args": {} } }CassandraTableSensor:等待表被创建
CassandraTableSensor用于探测 Cassandra 集群中某个表是否存在,完整实现见 sensors/table.py。它的构造函数签名如下:
CassandraTableSensor( *, table: str, cassandra_conn_id: str = "cassandra_default", **kwargs, )table:目标 Cassandra 表,使用点号(dot notation)可精确定位到指定 keyspace,例如keyspace_name.table_name;若不写 keyspace,则使用连接中 Schema 字段指定的 keyspace。cassandra_conn_id:连接 Cassandra 集群所用的连接 ID,默认cassandra_default。
该传感器继承自BaseSensorOperator,并把table声明为template_fields(源码见 table.py#L52),意味着table支持 Jinja 模板渲染,可动态传入上游产出的表名。
其核心探测逻辑poke()非常简单(源码见 table.py#L65-L68):
def poke(self, context: Context) -> bool: self.log.info("Sensor check existence of table: %s", self.table) hook = CassandraHook(self.cassandra_conn_id) return hook.table_exists(self.table)每一次 poke 都会实例化CassandraHook并调用hook.table_exists(self.table):返回True表示表已存在,传感器立即成功;返回False则继续等待,直到BaseSensorOperator配置的poke_interval(默认 60 秒)与timeout(默认 7 天)等超时参数生效。因此在使用时,建议同时设置合理的poke_interval与timeout,例如poke_interval=30, timeout=600,避免探测过于频繁或无限期等待。
底层实现:table_exists 是如何判定的
CassandraHook.table_exists()的实现见 hooks/cassandra.py#L176-L187:
def table_exists(self, table: str) -> bool: keyspace = self.keyspace if "." in table: keyspace, table = table.split(".", 1) cluster_metadata = self.get_conn().cluster.metadata return keyspace in cluster_metadata.keyspaces and table in cluster_metadata.keyspaces[keyspace].tables可以看到它并不执行 CQL 查询,而是直接读取Cluster.metadata(集群元数据缓存):先从表名中解析出 keyspace(若包含点号则拆分,否则使用连接 Schema),再检查该 keyspace 是否存在于元数据中、且目标表是否在该 keyspace 的 tables 集合中。这种基于元数据的方式开销极小,非常适合高频轮询。
集成测试 test_cassandra.py#L205-L235 验证了两种定位方式:hook.table_exists("s.t")从 CQL 字符串解析 keyspace,hook.table_exists("t")则使用 session 的 keyspace,二者均能正确返回 True/False。
CassandraRecordSensor:等待记录被写入
CassandraRecordSensor用于探测某个表中是否存在满足条件的记录,完整实现见 sensors/record.py。构造函数签名如下:
CassandraRecordSensor( *, keys: dict[str, str], table: str, cassandra_conn_id: str = "cassandra_default", **kwargs, )table:目标表,使用点号定位 keyspace,例如keyspace_name.table_name。keys:需要探测的键值对字典,例如{"p1": "v1", "p2": "v2"},表示“等待列p1的值为v1且列p2的值为v2的记录出现”。该字典的每个键都会作为 CQL WHERE 条件,多个条件用 AND 连接。cassandra_conn_id:连接 ID,默认cassandra_default。
该传感器同样继承自BaseSensorOperator,且template_fields同时包含("table", "keys")(源码见 record.py#L56),table与keys都支持 Jinja 模板。
poke()实现见 record.py#L71-L74:
def poke(self, context: Context) -> bool: self.log.info("Sensor check existence of record: %s", self.keys) hook = CassandraHook(self.cassandra_conn_id) return hook.record_exists(self.table, self.keys)底层实现:record_exists 与 CQL 构建
CassandraHook.record_exists()实现见 hooks/cassandra.py#L195-L214:
def record_exists(self, table: str, keys: dict[str, str]) -> bool: keyspace = self._sanitize_input(self.keyspace) if self.keyspace else self.keyspace if "." in table: keyspace, table = map(self._sanitize_input, table.split(".", 1)) else: table = self._sanitize_input(table) ks_str = " AND ".join(f"{key}=%({key})s" for key in keys) query = f"SELECT * FROM {keyspace}.{table} WHERE {ks_str}" try: result = self.get_conn().execute(query, keys) return result.one() is not None except Exception: return False该实现有两点值得注意:
- 输入消毒(防止 CQL 注入):
_sanitize_input()(hooks/cassandra.py#L189-L193)使用正则^\w+$校验 keyspace 与表名,只允许字母、数字、下划线,否则抛出ValueError。集成测试 test_cassandra.py#L237-L251 专门验证了record_exists("t; DROP TABLE t; SELECT * FROM t", ...)会被拒绝并抛出Invalid input异常。 - 参数化绑定:
keys字典以命名参数(%(key)s)形式绑定到 CQL 查询中,键值本身不会拼接进语句,进一步杜绝注入风险。 - 异常兜底:查询任何异常(如表不存在、查询超时)都会返回
False,传感器继续等待,而不会直接让任务失败——这与“等待就绪”的语义是一致的。
同时注意,record_exists使用result.one() is not None判断记录是否存在,只取第一条结果,因此适用于按主键精确匹配的场景。
完整示例:在 DAG 中使用两个传感器
官方在 tests/system/apache/cassandra/example_cassandra_dag.py 中提供了一个可直接参考的示例 DAG,其中用default_args集中声明了table,并用点号定位 keyspace(keyspace_name.table_name),完整代码如下:
from __future__ import annotations import os from datetime import datetime from airflow.models import DAG from airflow.providers.apache.cassandra.sensors.record import CassandraRecordSensor from airflow.providers.apache.cassandra.sensors.table import CassandraTableSensor ENV_ID = os.environ.get("SYSTEM_TESTS_ENV_ID") DAG_ID = "example_cassandra_operator" with DAG( dag_id=DAG_ID, 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放进default_args后,两个传感器无需重复书写表名;单任务内也可以用table="keyspace.table"显式覆盖。 CassandraRecordSensor的keys={"p1": "v1", "p2": "v2"}表示探测“列p1值为v1、列p2值为v2的记录”。- 实际部署时请把
<table_name>替换为真实表名,例如analytics.events;该示例文件同时可作为系统测试运行(文件末尾通过get_test_run(dag)接入 pytest 系统测试框架,运行方式见 系统测试文档)。 - 若需要串行“先等表、再等记录”的依赖关系,可在 DAG 中显式声明
record_sensor >> table_sensor或反之,按业务语义编排顺序。
传感器参数继承:BaseSensorOperator 的通用能力
由于CassandraTableSensor与CassandraRecordSensor均继承自 Airflow 的BaseSensorOperator(本仓库中通过airflow.providers.common.compat.sdk兼容层引入,见 table.py#L24 与 record.py#L24),它们自动拥有传感器家族的通用参数,可据此控制等待行为:
| 参数 | 默认值 | 说明 |
|---|---|---|
poke_interval | 60(秒) | 两次探测之间的间隔 |
timeout | 604800(7 天,秒) | 超过该时长仍未满足条件则任务失败 |
mode | poke | poke模式占用一个 worker 槽位持续轮询;reschedule模式在两次探测之间释放槽位,更适合长时间等待 |
exponential_backoff | False | 开启后探测间隔按指数退避增长 |
max_retry_delay | 1.0(天,秒) | 指数退避时单次最大间隔 |
针对 Cassandra 的“等数据就绪”场景,推荐组合使用mode="reschedule"、合理的poke_interval(如 30 秒)与显式timeout,在节省调度资源的同时保证任务最终超时可控。
深入底层:CassandraHook 如何构建连接
两个传感器的poke()都通过CassandraHook与集群通信,理解 Hook 的连接构建过程有助于排查“传感器一直不满足”或“连接失败”类问题。CassandraHook.__init__(hooks/cassandra.py#L89-L123)会将 Connection 字段映射为cassandra-driver的Cluster配置:
conn.host:按逗号拆分得到contact_points(支持多节点列表);conn.port:转换为port;conn.login/conn.password:构建PlainTextAuthProvider用于认证;conn.schema:作为默认 keyspace 保存到self.keyspace,get_conn()建立 session 时使用;- Extra 字段依次解析
load_balancing_policy、load_balancing_policy_args、cql_version、ssl_options、protocol_version并传入Cluster。
其中负载均衡策略由get_lb_policy()(hooks/cassandra.py#L141-L174)统一工厂化创建,行为如下:
DCAwareRoundRobinPolicy:读取local_dc(本地数据中心,默认空串)与used_hosts_per_remote_dc(每个远端数据中心使用的主机数,默认 0);WhiteListRoundRobinPolicy:必须提供hosts列表,否则抛出ValueError;TokenAwarePolicy:默认子策略为RoundRobinPolicy,也可通过child_load_balancing_policy指定RoundRobinPolicy/DCAwareRoundRobinPolicy/WhiteListRoundRobinPolicy之一;- 其他任何名称(含非法策略名)都会回退到默认的
RoundRobinPolicy。
这些行为在集成测试 test_cassandra.py#L88-L151 中都有对应断言:非法策略名回退RoundRobinPolicy、TokenAware 子策略默认RoundRobinPolicy、白名单策略缺少 hosts 抛异常等。需要注意的是,官方连接文档中的AllowListRoundRobinPolicy命名与源码中实际实现的WhiteListRoundRobinPolicy(见 hooks/cassandra.py#L154-L158)存在新旧命名差异,配置 Extra 时以源码实际支持的策略名为准。
此外,Hook 提供get_cluster()、shutdown_cluster()等方法管理集群生命周期,get_conn()会缓存并复用Session(hooks/cassandra.py#L125-L130),避免每次探测重复建连。
测试验证:单元测试与集成测试一览
仓库为该 Provider 提供了完整的测试覆盖,可作为理解传感器契约的补充材料:
- sensors/test_table.py:通过 mock
CassandraHook验证CassandraTableSensor.poke()会以正确的连接 ID 与表名调用hook.table_exists;并验证不传cassandra_conn_id时默认值为cassandra_default,以及带 keyspace 点号表名的场景。 - sensors/test_record.py:验证
CassandraRecordSensor.poke()以(table, keys)调用hook.record_exists,包括keys=None的边界场景与默认连接 ID。 - hooks/test_cassandra.py:标记为
@pytest.mark.integration("cassandra"),需要真实 Cassandra 环境,覆盖连接构建(多 host、端口、协议版本、负载均衡策略)、四种负载均衡策略的参数化创建、table_exists/record_exists在“字符串含 keyspace”与“session 默认 keyspace”两种模式下的正确性,以及 CQL 注入防护。
实战注意事项与最佳实践
- 表名务必使用点号或正确配置 Schema:若连接 Schema 未设置且表名不带 keyspace,
record_exists会构造出FROM .table这类非法语句并因异常返回False,传感器将永远等待直到超时。 - 为传感器设置合理的超时:默认
timeout长达 7 天,生产环境建议显式配置timeout与poke_interval,避免资源被长期占用;长等待场景优先使用mode="reschedule"。 - 探测键值应使用主键:
record_exists生成的SELECT ... WHERE p1=v1 AND p2=v2在非主键列上需要全表扫描,性能较差;请尽量以分区键/主键作为keys条件。 - 多数据中心部署时配置负载均衡策略:通过 Connection 的 Extra 字段指定
DCAwareRoundRobinPolicy并设置local_dc,可显著降低跨数据中心延迟;启用 SSL 的集群记得配置ssl_options。 - 利用模板字段实现动态表名:
table与keys都是template_fields,可从上游任务产出的 XCom 或调度参数动态渲染,例如table="analytics.{{ ds_nodash }}"按日期分表等待。
综上,CassandraTableSensor与CassandraRecordSensor提供了“等待表就绪”与“等待数据就绪”两种轻量级探测能力,配合CassandraHook的元数据探测、参数化查询与注入防护机制,能够让 Airflow DAG 与 Cassandra 数据管道实现可靠、安全的时序对齐。
【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考