如何通过PySpark在Databricks中使用私钥连接SFTP直接读取文件?
在Databricks中用PySpark通过私钥直接读取SFTP文件的方法
有两种主流方法可以实现Databricks直接读取SFTP文件,无需中转到Linux服务器或Azure容器:
方法1:使用Spark SFTP数据源(推荐,适合批量/大文件)
这个方法依赖Spark生态的SFTP数据源,性能更优,适合直接读取结构化文件。
安装依赖
在Databricks集群中添加Maven库依赖:com.springml:spark-sftp_2.12:1.1.1(Scala版本需匹配你的Databricks集群版本,比如2.12对应Spark 3.x)。也可以在Notebook中执行%pip install spark-sftp安装Python端依赖。配置私钥与读取文件
将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。
- 安装Paramiko
在Notebook中执行:
%pip install paramiko
- 代码示例
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
相关产品推荐
相关产品推荐

