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

如何在PySpark中统计指定时间戳后目录及子目录内的Avro文件数量并存储到变量

Got it, let's break down how to solve this problem in PySpark. Since we need to interact directly with the filesystem to check file timestamps and count specific files, we'll combine PySpark with Hadoop's FileSystem API (which Spark relies on under the hood) to get the job done.

Solution Steps & Code Example

Core Approach

  • Use Spark's context to access the Hadoop FileSystem instance, which lets us traverse directories and fetch file metadata
  • Filter for files with the .avro extension
  • Compare each file's modification time (note: HDFS doesn't natively track "creation time"—we typically use modification time as a stand-in; adjust if your filesystem supports creation time) against your target timestamp
  • Recursively scan subdirectories and count qualifying files, storing the result in a variable

Full Code Implementation

from pyspark.sql import SparkSession
import time

def count_avro_files_after_timestamp(spark, target_dir, target_timestamp):
    # Get Hadoop configuration from Spark
    hadoop_conf = spark.sparkContext._jsc.hadoopConfiguration()
    # Initialize Hadoop FileSystem instance
    fs = spark.sparkContext._jvm.org.apache.hadoop.fs.FileSystem.get(hadoop_conf)
    # Convert target directory path to Hadoop's Path object
    hdfs_path = spark.sparkContext._jvm.org.apache.hadoop.fs.Path(target_dir)
    
    file_count = 0
    # List all file/directory statuses in the target path
    for status in fs.listStatus(hdfs_path):
        if status.isDirectory():
            # Recursively count files in subdirectories
            file_count += count_avro_files_after_timestamp(
                spark, status.getPath().toString(), target_timestamp
            )
        else:
            # Check if the file is an Avro file
            file_path = status.getPath().toString()
            if file_path.endswith(".avro"):
                # Get file modification time (in milliseconds)
                file_modify_ms = status.getModificationTime()
                # Convert target timestamp to milliseconds (assuming input is in seconds)
                target_ms = target_timestamp * 1000
                if file_modify_ms >= target_ms:
                    file_count += 1
    return file_count

# Initialize SparkSession
spark = SparkSession.builder.appName("AvroFileCounter").getOrCreate()

# Configuration parameters
target_directory = "/path/to/your/target/directory"
# Example: Timestamp for Jan 1, 2024 (UTC, in seconds)
cutoff_timestamp = 1704067200

# Execute the count
qualified_avro_count = count_avro_files_after_timestamp(spark, target_directory, cutoff_timestamp)
print(f"Total Avro files created/modified after the cutoff: {qualified_avro_count}")

# Clean up SparkSession
spark.stop()

Key Details to Note

  1. Hadoop API Interaction: We use Spark's _jvm attribute to access Java Hadoop classes since PySpark doesn't have a native Python API for filesystem metadata operations. This works across both local and cluster environments (like YARN + HDFS).
  2. Timestamp Conversion: Hadoop returns file timestamps in milliseconds, so if your input timestamp is in seconds (common in Unix timestamps), multiply by 1000 to match.
  3. Creation Time vs Modification Time: HDFS doesn't track creation time natively. If you're working with a local filesystem (e.g., ext4) and need actual creation time, you might need to use Python's os module instead—but note this won't work in a distributed cluster where files are on remote nodes.
  4. Performance for Large Directories: If you're dealing with tens of thousands of files, the recursive approach above might be slow. For distributed scaling, try this alternative using Spark's built-in file handling:
from pyspark.sql.functions import input_file_name, udf
from pyspark.sql.types import LongType

def get_modify_time(path):
    hadoop_conf = spark.sparkContext._jsc.hadoopConfiguration()
    fs = spark.sparkContext._jvm.org.apache.hadoop.fs.FileSystem.get(hadoop_conf)
    return fs.getFileStatus(spark.sparkContext._jvm.org.apache.hadoop.fs.Path(path)).getModificationTime()

# Create UDF to fetch modification time
get_modify_time_udf = udf(get_modify_time, LongType())

# Load Avro file paths (no need to load full data)
file_paths_df = spark.read.format("avro")\
    .option("recursiveFileLookup", "true")\
    .load(target_directory)\
    .select(input_file_name().alias("file_path"))\
    .distinct()  # Avoid counting the same file multiple times

# Add modification time column and filter
target_ms = cutoff_timestamp * 1000
qualified_count = file_paths_df\
    .withColumn("modify_time", get_modify_time_udf("file_path"))\
    .filter(f"modify_time >= {target_ms}")\
    .count()

This leverages Spark's distributed processing to handle large file sets more efficiently.

内容的提问来源于stack exchange,提问作者bda

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.29 00:52:45