【问题标题】:How to insert record in cassandra table using airflow?如何使用气流在 cassandra 表中插入记录?
【发布时间】:2021-07-26 18:44:00
【问题描述】:

我已经将 Cassandra 和气流一起安装在 docker 中。

我想使用气流在 Cassandra 表中插入数据。

就像airflow有MySqlOperator在SQL表中安装数据一样,有没有什么操作符或者方法可以将记录插入到Cassandra表中。

我只找到了这两个运算符: 从airflow.providers.apache.cassandra.sensors.record 导入CassandraRecordSensor 从airflow.providers.apache.cassandra.sensors.table 导入CassandraTableSensor

但这些运算符只是为了检查表或记录cassandra中的存在。

那么,如何使用气流任务插入或说与 Cassandra 交互?

【问题讨论】:

    标签: docker cassandra airflow


    【解决方案1】:

    文档显示确实没有实现“写入”操作:

    https://airflow.apache.org/docs/apache-airflow-providers-apache-cassandra/stable/operators.html

    但是如果你没有现成的操作符,Apache Airflow 真的很容易扩展。

    如果您了解 Python 方法,则需要扩展 Cassandra Hook 并实现自定义运算符(并可能在您这样做时将其回馈给社区)。这是最好的,因为您将能够使用我猜已经存在的 cassandra 库和身份验证。

    或者您可以使用 BashOperator 来运行 CQL 命令(我相信这是 cassandra 使用的默认客户端)。例如,如果您有 CSV 文件,您可以使用 CQL 中的 COPY 命令将其导入。

    https://docs.datastax.com/en/cql-oss/3.x/cql/cql_reference/cqlshCopy.html

    然后您必须在来自连接的身份验证信息之间进行一些链接,并将其传递给 BashOperator,或者提供您自己的方式来使用 Cassandra 进行身份验证。

    【讨论】:

    • 我尝试了您的第一个解决方案,即创建自定义 Cassandra 运算符。唯一的问题是我无法导入 CassandraHook,它给出了一个错误。我尝试了以下 2 种方法: from airflow.providers.apache.cassandra.hooks.cassandra import CassandraHook from airflow.contrib.hooks.cassandra_hook import CassandraHook 扩展 CassandraHook 的正确方法是什么?
    • 如果您使用的是 Airlfow 1.10 - 升级到 Airflow 2。 Airflow 1.10 已于 6 月结束生命周期,它不会再得到任何修复 - 甚至是关键的安全修复。所以尽快升级到Airflow 2。如果你在 2 上,你需要安装 'apache-airflow-providers-apache-cassandra' 来获取这些包。 Airflow 2 已被拆分为“核心”气流和约 70 个供应商。
    猜你喜欢
    • 2021-01-31
    • 2017-09-15
    • 1970-01-01
    • 2019-11-21
    • 2017-04-19
    • 1970-01-01
    • 1970-01-01
    • 2018-12-03
    • 1970-01-01
    相关资源
    最近更新 更多