Vane Data / 核心概念
SQL 与 Python
SQL 和 Python 是构建同一条惰性 Vane Data pipeline 的两种接口。你可以 单独使用任一接口,也可以在一条 pipeline 中混合使用:两者都会生成 Relation 对象,并且任一接口都可以继续消费这些对象。
SQL 接口
当转换逻辑最适合用 SQL 语句表达时,使用 con.sql(...)。source 等 Python Relation 变量可以直接按名称引用:
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 对象:
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 聚合结果:
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()