Rust DataFusion DataFrame从Unix时间戳列提取年月日方法
实现方案
DataFusion 内置了和 PySpark 对齐的日期时间处理函数,不需要自定义 UDF 即可完成需求,实现逻辑和你写的 PySpark 代码几乎一致。
核心函数说明
- 时间戳转换:使用内置
from_unixtime函数,接收秒级 Unix 时间戳列,返回 DataFusion 原生Timestamp类型列,可直接和 Rust chrono 等时间库的类型交互。 - 日期字段提取:使用内置
year、month、day函数,传入时间戳列即可返回对应整数类型的年、月、日值。
完整实现代码
首先在Cargo.toml中引入对应版本的 DataFusion 依赖:
[dependencies] datafusion = "38.0.0" tokio = { version = "1.0", features = ["rt-multi-thread", "macros"] }
Rust 实现代码:
use datafusion::prelude::*; use datafusion::functions::datetime::{day, from_unixtime, month, year}; #[tokio::main] async fn main() -> datafusion::error::Result<()> { let ctx = SessionContext::new(); // 读取你的JSON格式源数据 let main_df = ctx .read_json("你的数据文件路径.json", NdJsonReadOptions::default()) .await?; // 转换unix时间戳为原生时间戳类型 let main_df = main_df.with_column( "reviewed_at", from_unixtime(vec![col("unixReviewTime")])?, )?; // 分别提取年、月、日为独立列 let main_df = main_df .with_column("reviewed_year", year(col("reviewed_at"))?)? .with_column("reviewed_month", month(col("reviewed_at"))?)? .with_column("reviewed_day", day(col("reviewed_at"))?)?; // 打印结果验证 main_df.show_limit(10).await?; Ok(()) }
注意事项
- 上述代码默认
unixReviewTime是秒级时间戳,和你给出的样例数据匹配。如果是毫秒级时间戳,替换from_unixtime为to_timestamp_millis函数即可,不需要做除法转换。 - 如果需要指定时区提取日期,可以给
from_unixtime传入第二个时区参数,例如from_unixtime(vec![col("unixReviewTime"), lit("Asia/Shanghai")])?,后续提取年月日时会按照指定时区计算,默认是UTC时区。 - 转换完成后通过
schema()方法查看表结构,可确认reviewed_at为Timestamp类型,年/月/日列为Int32类型,符合Rust原生类型交互要求。
内容的提问来源于stack exchange,提问作者DataPsycho
相关产品推荐
相关产品推荐

