
Apache Airflow 集成 Apache CassandraTableSensor 与 RecordSensor 实战指南与源码解析【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflowApache 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表中是否存在tableCassandraRecordSensor表中是否存在指定记录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默认RoundRobinPolicyload_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-L68def 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_interval30, timeout600避免探测过于频繁或无限期等待。底层实现table_exists 是如何判定的CassandraHook.table_exists()的实现见 hooks/cassandra.py#L176-L187def 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 字符串解析 keyspacehook.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#L56table与keys都支持 Jinja 模板。poke()实现见 record.py#L71-L74def 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-L214def 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 fSELECT * 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并用点号定位 keyspacekeyspace_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_idDAG_ID, scheduleNone, start_datedatetime(2021, 1, 1), default_args{table: keyspace_name.table_name}, catchupFalse, tags[example], ) as dag: # Replace table_name with your actual table name table_sensor CassandraTableSensor(task_idcassandra_table_sensor, tabletable_name) record_sensor CassandraRecordSensor( task_idcassandra_record_sensor, keys{p1: v1, p2: v2}, tabletable_name )要点说明把table放进default_args后两个传感器无需重复书写表名单任务内也可以用tablekeyspace.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_interval60秒两次探测之间的间隔timeout6048007 天秒超过该时长仍未满足条件则任务失败modepokepoke模式占用一个 worker 槽位持续轮询reschedule模式在两次探测之间释放槽位更适合长时间等待exponential_backoffFalse开启后探测间隔按指数退避增长max_retry_delay1.0天秒指数退避时单次最大间隔针对 Cassandra 的“等数据就绪”场景推荐组合使用modereschedule、合理的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转换为portconn.login/conn.password构建PlainTextAuthProvider用于认证conn.schema作为默认 keyspace 保存到self.keyspaceget_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每个远端数据中心使用的主机数默认 0WhiteListRoundRobinPolicy必须提供hosts列表否则抛出ValueErrorTokenAwarePolicy默认子策略为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()会缓存并复用Sessionhooks/cassandra.py#L125-L130避免每次探测重复建连。测试验证单元测试与集成测试一览仓库为该 Provider 提供了完整的测试覆盖可作为理解传感器契约的补充材料sensors/test_table.py通过 mockCassandraHook验证CassandraTableSensor.poke()会以正确的连接 ID 与表名调用hook.table_exists并验证不传cassandra_conn_id时默认值为cassandra_default以及带 keyspace 点号表名的场景。sensors/test_record.py验证CassandraRecordSensor.poke()以(table, keys)调用hook.record_exists包括keysNone的边界场景与默认连接 ID。hooks/test_cassandra.py标记为pytest.mark.integration(cassandra)需要真实 Cassandra 环境覆盖连接构建多 host、端口、协议版本、负载均衡策略、四种负载均衡策略的参数化创建、table_exists/record_exists在“字符串含 keyspace”与“session 默认 keyspace”两种模式下的正确性以及 CQL 注入防护。实战注意事项与最佳实践表名务必使用点号或正确配置 Schema若连接 Schema 未设置且表名不带 keyspacerecord_exists会构造出FROM .table这类非法语句并因异常返回False传感器将永远等待直到超时。为传感器设置合理的超时默认timeout长达 7 天生产环境建议显式配置timeout与poke_interval避免资源被长期占用长等待场景优先使用modereschedule。探测键值应使用主键record_exists生成的SELECT ... WHERE p1v1 AND p2v2在非主键列上需要全表扫描性能较差请尽量以分区键/主键作为keys条件。多数据中心部署时配置负载均衡策略通过 Connection 的 Extra 字段指定DCAwareRoundRobinPolicy并设置local_dc可显著降低跨数据中心延迟启用 SSL 的集群记得配置ssl_options。利用模板字段实现动态表名table与keys都是template_fields可从上游任务产出的 XCom 或调度参数动态渲染例如tableanalytics.{{ 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),仅供参考