ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

Apache Airflow 之 Apache Livy 连接配置详解:从 Connection 字段到 LivyHook 底层实现

Apache Airflow 之 Apache Livy 连接配置详解:从 Connection 字段到 LivyHook 底层实现 Apache Airflow 之 Apache Livy 连接配置详解从 Connection 字段到 LivyHook 底层实现【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow在 Apache Airflow 中对接 Spark 集群通常通过 Apache Livy 的 REST API 完成Livy 服务以简单的 HTTP 接口暴露 Spark 批处理、会话管理与结果查询能力。Airflow 的 Apache Livy Provider 为此定义了专用的livy连接类型Connection Type并通过LivyHook消费该连接。本文基于 Provider 官方文档 connections.rst 展开结合仓库源码与测试用例讲清楚Livy Connection 的每个字段Host、Port、Schema、Login、Password、Extra如何被底层解析为请求 URL 与认证信息以及如何通过环境变量以 URI 语法配置连接帮助读者完成一个可复制、可运行的 Livy 连接配置并理解其背后的实现机制。什么是 Livy 连接类型Airflow 对 Livy 的集成入口是LivyHook其实现位于 hooks/livy.py。从源码定义可以直接看到它与 HTTP 连接体系的从属关系class LivyHook(HttpHook): Hook for Apache Livy through the REST API. conn_name_attr livy_conn_id default_conn_name livy_default conn_type livy hook_name Apache Livy default_headers {Content-Type: application/json, Accept: application/json}也就是说Apache Livy 连接在底层就是 HTTP 连接LivyHook直接继承自 Providers 中的HttpHookhttp.py复用了 HTTP 连接的全部 URL 构建、认证与重试机制只在此基础上封装了 Livy REST API 的端点调用。Provider 的元数据文件 provider.yaml 中显式注册了该映射connection-types: - hook-class-name: airflow.providers.apache.livy.hooks.livy.LivyHook hook-name: Apache Livy connection-type: livy默认连接 IDlivy_default官方文档指出Livy Hook 使用参数livy_conn_id作为 Connection ID其默认值为livy_default。这一点在源码中同样得到印证——default_conn_name livy_default是类属性因此当你在 DAG 中实例化LivyHook()或LivyOperator等组件时不显式指定连接 ID就会自动回退到名为livy_default的连接。这一设计带来两个实用结论如果整条流水线只用一个 Livy 服务直接创建livy_default连接即可所有组件零配置如果集群中有多个 Livy 服务例如不同队列或不同版本则创建多个连接如livy_staging、livy_prod并在组件处显式传入livy_conn_idlivy_staging覆盖默认值。配置字段逐项解析在 Airflow UI 的 Connections 页面或airflow connections add中配置 Livy 连接时文档定义了以下字段下面逐项说明其含义与底层解析行为。HostLivy 服务器的 HTTP 主机名。文档允许两种写法在Host字段中直接写入带 scheme 的完整地址例如http://livy-server.com只填主机名把协议放到Schema字段中。源码层面HttpHook的_set_base_url方法http.py完整实现了这一约定def _set_base_url(self, connection) - None: host connection.host or self.default_host schema connection.schema or http # RFC 3986 (https://www.rfc-editor.org/rfc/rfc3986.html#page-16) if :// in host: self.base_url host else: self.base_url f{schema}://{host} if host else f{schema}:// if connection.port: self.base_url f{self.base_url}:{connection.port} parsed urlparse(self.base_url) if not parsed.scheme: raise ValueError(fInvalid base URL: Missing scheme in {self.base_url})从中可以确认几条关键规则若Host中已含://即自带 scheme则Host整体作为 base URLSchema字段被忽略——这避免了https://被二次拼接若Host是裸主机名则拼成{schema}://{host}schema缺省为http最终 base URL 必须能被urlparse解析出 scheme否则抛出ValueError。Port如果 Host 是完整的 URL 形式端口可以不填否则通过Port字段指定最终被追加到 base URL 尾部如http://livy-server.com:8998。Livy 的默认服务端口通常是 8998仓库单元测试中同样以port8998作为测试连接的默认端口test_livy.py。Schema可选指定协议类型取值为http或https。注意优先级问题只要Host里写了 schemeSchema就不生效源码中:: in host分支直接返回 host。这一“Host 自带 scheme 优先”的行为在单元测试的参数化用例中被明确验证见 test_livy.pypytest.mark.parametrize( (conn_id, expected), [ pytest.param(default_port, http://host, iddefault-port), pytest.param(default_protocol, http://host, iddefault-protocol), pytest.param(port_set, http://host:1234, idwith-defined-port), pytest.param(schema_set, https://host, idwith-defined-schema), pytest.param(dont_override_schema, http://host, idignore-defined-schema), ], ) def test_build_get_hook(self, conn_id, expected): hook LivyHook(livy_conn_idconn_id) hook.get_conn() assert hook.base_url expected其中ignore-defined-schema用例专门覆盖“Host 已含http://时schemahttps不生效”的场景与文档描述的优先级一致。Login 与 Password可选分别指定 Livy 服务器的登录名与密码。底层由HttpHook的认证提取逻辑处理http.py只要connection.login非空就以login/password构造认证凭据挂载到 requests Session 上。Livy 集群若启用了 HTTP Basic Auth 或经过反向代理鉴权就需要填写这两个字段。Extra可选以 JSON 格式指定附加 HTTP 请求头。文档给出的语义是 “Specify headers in json format”即 Extra 字段的值是包含headers键的 JSON。HttpHook在构建 Session 时通过_configure_session_from_extra解析 Connection 的extra字段把其中的 headers 合并进请求头http.pyif connection.extra or extra_options: session self._configure_session_from_extra(session, connection, extra_options) ... if self.default_headers: session.headers.update(self.default_headers)注意请求头的合并顺序LivyHook的default_headersContent-Type: application/json与Accept: application/json会在 Session 创建后统一更新确保 Livy REST API 所需的 JSON 语义始终成立。通过环境变量配置连接如果偏好环境变量而非 UI文档要求使用 URI 语法并强调URI 的所有组成部分都应做 URL 编码。官方示例export AIRFLOW_CONN_LIVY_DEFAULTlivy://username:passwordlivy-server.com:80/http?headersheader环境变量名遵循AIRFLOW_CONN_CONN_ID 的大写形式规则AIRFLOW_CONN_LIVY_DEFAULT即对应默认连接 IDlivy_default。URI 中各段与 Connection 字段的映射关系为URI 段对应字段示例schemelivy://连接类型livyusername:passwordLogin / Password认证凭据livy-server.comHost服务器地址:80Port端口查询参数Extra附加信息如 headers配置完成后建议先用轻量方式验证连通性——LivyHook.run_method支持GET方法可直接请求 Livy 的/batches等只读端点确认网络与认证是否打通。连接在运行期如何被使用配置好的连接由LivyHook消费围绕 Livy REST API 封装了完整的批处理生命周期调用链hooks/livy.pypost_batch(...)以POST {prefix}/batches提交 Spark 批作业请求体由静态方法build_post_batch_body构建支持file、class_name、args、jars、py_files、queue、driver_memory、num_executors、conf等 Livy 批接口字段成功后解析响应中的id返回 batch session id。get_batch(session_id)/get_batch_state(session_id)分别GET {prefix}/batches/{id}与{prefix}/batches/{id}/state后者把响应中的 state 映射为BatchState枚举NOT_STARTED、STARTING、RUNNING、IDLE、BUSY、SHUTTING_DOWN、ERROR、DEAD、KILLED、SUCCESS其中SUCCESS/DEAD/KILLED/ERROR被定义为TERMINAL_STATES供轮询类组件判断作业是否结束。delete_batch(session_id)DELETE {prefix}/batches/{id}终止会话。get_batch_logs(...)/dump_batch_logs(...)分页拉取GET {prefix}/batches/{id}/logdump_batch_logs以每批 100 行的方式循环抓取并写入任务日志。run_method(...)统一的请求入口仅允许GET/POST/PUT/DELETE/HEAD并可通过retry_args走run_with_advanced_retry高级重试。此外还提供LivyAsyncHook继承自HttpAsyncHook用于 Airflow 的 deferrable 模式LivyOperator设置deferrableTrue后状态轮询交给 Triggerer 异步完成减少 worker 占用对应的触发器实现在 triggers/livy.py。一个真实的使用示例见系统测试 DAG example_livy.py它直接使用默认连接livy_default提交 Java 与 Python 两种 Spark 计算livy_java_task LivyOperator( task_idpi_java_task, file/spark-examples.jar, num_executors1, conf{spark.shuffle.compress: false}, class_nameorg.apache.spark.examples.SparkPi, ) livy_python_task LivyOperator(task_idpi_python_task, file/pi.py, polling_interval60)更多 Operator 用法可参考 operators.rst。小结与实践建议Livy 连接的本质是 HTTP 连接Host 可自带 schemeSchema 只作为裸主机名时的协议补充且自带 scheme 时优先级更高端口、认证、附加 headers 均由HttpHook统一处理。默认连接 ID 固定为livy_defaultdefault_conn_name多集群场景通过livy_conn_id参数切换。环境变量配置必须使用livy://URI 语法并对各组件做 URL 编码变量名与连接 ID 保持大小写映射AIRFLOW_CONN_LIVY_DEFAULT↔livy_default。配置完成后LivyHook即可以Content-Type: application/json语义调用 Livy 批处理 API长作业建议搭配deferrableTrue由 Triggerer 异步轮询。【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表