跳到主要内容
Vane Data / 核心概念

SQL 与 Python

SQL 和 Python 是构建同一条惰性 Vane Data pipeline 的两种接口。你可以 单独使用任一接口,也可以在一条 pipeline 中混合使用:两者都会生成 Relation 对象,并且任一接口都可以继续消费这些对象。

SQL 接口

当转换逻辑最适合用 SQL 语句表达时,使用 con.sql(...)source 等 Python Relation 变量可以直接按名称引用:

example.py
import vane


con = vane.connect()
source = con.values(
    (
        vane.lit(1).alias("order_id"),
        vane.lit("books").alias("category"),
        vane.lit(2).alias("quantity"),
        vane.lit(12).alias("unit_price"),
    ),
    (vane.lit(2), vane.lit("games"), vane.lit(1), vane.lit(60)),
    (vane.lit(3), vane.lit("books"), vane.lit(3), vane.lit(8)),
    (vane.lit(4), vane.lit("stationery"), vane.lit(1), vane.lit(5)),
)


result = con.sql("""
    SELECT
        order_id,
        category,
        quantity * unit_price AS total
    FROM source
    WHERE quantity * unit_price >= 20
    ORDER BY order_id
""")
result.show()

Python 接口

表操作使用 Relation 方法;列、字面量、算术运算和谓词使用 Expression 对象:

example.py
import vane


con = vane.connect()
source = con.values(
    (
        vane.lit(1).alias("order_id"),
        vane.lit("books").alias("category"),
        vane.lit(2).alias("quantity"),
        vane.lit(12).alias("unit_price"),
    ),
    (vane.lit(2), vane.lit("games"), vane.lit(1), vane.lit(60)),
    (vane.lit(3), vane.lit("books"), vane.lit(3), vane.lit(8)),
    (vane.lit(4), vane.lit("stationery"), vane.lit(1), vane.lit(5)),
)


result = source.filter(
    vane.col("quantity") * vane.col("unit_price") >= vane.lit(20)
).select(
    vane.col("order_id"),
    vane.col("category"),
    (
        vane.col("quantity") * vane.col("unit_price")
    ).alias("total"),
).order("order_id")
result.show()

混合 Pipeline

Relation 可以双向跨越接口边界。下面的示例先用 SQL 计算总价,再用 Python 筛选行,最后用 SQL 聚合结果:

example.py
import vane


con = vane.connect()
source = con.values(
    (
        vane.lit(1).alias("order_id"),
        vane.lit("books").alias("category"),
        vane.lit(2).alias("quantity"),
        vane.lit(12).alias("unit_price"),
    ),
    (vane.lit(2), vane.lit("games"), vane.lit(1), vane.lit(60)),
    (vane.lit(3), vane.lit("books"), vane.lit(3), vane.lit(8)),
    (vane.lit(4), vane.lit("stationery"), vane.lit(1), vane.lit(5)),
)


# SQL 阶段。
priced = con.sql("""
    SELECT
        order_id,
        category,
        quantity * unit_price AS total
    FROM source
""")


# Python 阶段。
selected = priced.filter(
    vane.col("total") >= vane.lit(20)
).select(
    vane.col("category"),
    vane.col("total"),
)


# SQL 阶段。
summary = con.sql("""
    SELECT category, CAST(sum(total) AS BIGINT) AS revenue
    FROM selected
    GROUP BY category
    ORDER BY category
""")
summary.show()