You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.27 15:54:19