【发布时间】:2020-11-06 17:42:29
【问题描述】:
我有一个 pyspark 管道,它应该将表格作为 CSV 文件导出到 HDFS 和 SFTP 服务器(之后数据将由 CRM 团队获取)。
要导出到 HDFS,它非常简单,而且效果很好, 但是要将数据导出到 sftp 文件,我这样做了:
def export_to_sftp():
dataframe.coalesce(1).options(codec=compression).write.mode("overwrite").option('encoding',encoding).csv(file_to_hdfs, header=True, nullValue='', sep=';')
copyToLocalFile(file_to_hdfs,local_machine_file) # copy from HDFS to LOCAL using hadoop API
cnopts = pysftp.CnOpts()
hostkeys = None
if cnopts.hostkeys.lookup(server) is not None:
hostkeys = cnopts.hostkeys
cnopts.hostkeys = None
try:
with pysftp.Connection(host=server, username=login,
password=password, cnopts=cnopts) as sftp:
if hostkeys is not None:
hostkeys.add(server, sftp.remote_server_key.get_name(), sftp.remote_server_key)
hostkeys.save(pysftp.helpers.known_hosts())
sftp.put(local_machine_file, sftp_path)
except Exception as e:
log.exception(e)
finally:
log.info("Cleaning up)
shutil.rmtree(local_tmp)
当文件不是太大时,此方法可以正常工作,但对于某些表它不起作用,因为在我的本地 linux 机器上我没有足够的磁盘空间,
那么是否可以使用 pysftp 以流的方式将远程文件从 HDFS 复制到 SFTP 而无需复制到本地机器?
【问题讨论】:
标签: python pyspark sftp pysftp