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

如何通过PySpark在Databricks中使用私钥连接SFTP直接读取文件?

在Databricks中用PySpark通过私钥直接读取SFTP文件的方法

有两种主流方法可以实现Databricks直接读取SFTP文件,无需中转到Linux服务器或Azure容器:

方法1:使用Spark SFTP数据源(推荐,适合批量/大文件)

这个方法依赖Spark生态的SFTP数据源,性能更优,适合直接读取结构化文件。

  1. 安装依赖
    在Databricks集群中添加Maven库依赖:com.springml:spark-sftp_2.12:1.1.1(Scala版本需匹配你的Databricks集群版本,比如2.12对应Spark 3.x)。也可以在Notebook中执行%pip install spark-sftp安装Python端依赖。

  2. 配置私钥与读取文件
    将SFTP私钥上传到DBFS(比如/dbfs/secrets/sftp_private_key.pem),然后用以下代码读取:

from pyspark.sql import SparkSession

spark = SparkSession.builder.appName("SFTPDirectRead").getOrCreate()

# 读取SFTP上的CSV文件示例
sftp_df = spark.read.format("com.springml.spark.sftp") \
    .option("host", "your-sftp-host.com") \
    .option("port", 22) \
    .option("username", "your-sftp-username") \
    .option("privateKey", "/dbfs/secrets/sftp_private_key.pem") \
    .option("privateKeyPassphrase", "your-key-passphrase")  # 私钥有密码时添加
    .option("fileType", "csv") \
    .option("header", "true")  # 文件带表头时添加
    .option("path", "/sftp-target-folder/*.csv") \
    .load()

# 查看数据
sftp_df.show()

注意事项:

  • 支持csv、parquet、json等多种文件类型,调整fileType参数即可
  • 路径支持通配符*匹配多个文件
  • 敏感信息(私钥、密码)建议用Databricks Secrets管理,避免硬编码:
    private_key_path = dbutils.secrets.get(scope="your-secret-scope", key="sftp-private-key-path")
    passphrase = dbutils.secrets.get(scope="your-secret-scope", key="sftp-key-passphrase")
    

方法2:使用Paramiko手动读取(适合自定义处理逻辑)

如果需要对文件内容做自定义预处理,可以用Paramiko库直接连接SFTP服务器,读取内容后转为Spark DataFrame。

  1. 安装Paramiko
    在Notebook中执行:
%pip install paramiko
  1. 代码示例
import paramiko
from pyspark.sql import Row
from io import StringIO

# 初始化SFTP客户端
ssh_client = paramiko.SSHClient()
ssh_client.set_missing_host_key_policy(paramiko.AutoAddPolicy())

# 加载私钥(从DBFS读取)
private_key = paramiko.RSAKey.from_private_key_file(
    "/dbfs/secrets/sftp_private_key.pem",
    password=dbutils.secrets.get(scope="your-secret-scope", key="sftp-key-passphrase")
)

# 连接SFTP服务器
ssh_client.connect(
    hostname="your-sftp-host.com",
    username="your-sftp-username",
    pkey=private_key
)
sftp = ssh_client.open_sftp()

# 遍历目标文件夹读取文件
target_folder = "/sftp-target-folder/"
data_rows = []

for file_name in sftp.listdir(target_folder):
    if file_name.endswith(".csv"):
        with sftp.open(f"{target_folder}{file_name}", "r") as f:
            # 读取文件内容并按行处理
            content = f.read().decode("utf-8")
            lines = content.split("\n")[1:]  # 跳过表头行
            for line in lines:
                if line.strip():
                    cols = line.split(",")
                    data_rows.append(Row(col1=cols[0], col2=cols[1], col3=cols[2]))  # 按实际列调整

# 转为Spark DataFrame
df = spark.createDataFrame(data_rows)
df.show()

# 关闭连接
sftp.close()
ssh_client.close()

注意事项:

  • 该方法适合小文件或需要自定义解析逻辑的场景,大文件建议用方法1
  • 务必处理文件读取后的连接关闭,避免资源泄漏

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.21 23:15:52