Delta: Error while writing data frame using s3a file system.

Created on 18 Jun 2019  Â·  3Comments  Â·  Source: delta-io/delta

Hi, I'm trying to use delta lake over s3 storage layer.

Spark command: df.write.format("delta").save("s3a://<bucket_name>/<key>")

Getting this error.

19/06/19 02:13:59 WARN DeltaLog: Failed to parse s3a://demo-atlan-source/delta_test/_delta_log/_last_checkpoint. This may happen if there was an error during read operation, or a file appears to be partial. Sleeping and trying again.
org.apache.hadoop.fs.UnsupportedFileSystemException: No AbstractFileSystem for scheme: s3a
    at org.apache.hadoop.fs.AbstractFileSystem.createFileSystem(AbstractFileSystem.java:154)
    at org.apache.hadoop.fs.AbstractFileSystem.get(AbstractFileSystem.java:242)
    at org.apache.hadoop.fs.FileContext$2.run(FileContext.java:334)
    at org.apache.hadoop.fs.FileContext$2.run(FileContext.java:331)
    at java.security.AccessController.doPrivileged(Native Method)
    at javax.security.auth.Subject.doAs(Subject.java:422)
    at org.apache.hadoop.security.UserGroupInformation.doAs(UserGroupInformation.java:1692)
    at org.apache.hadoop.fs.FileContext.getAbstractFileSystem(FileContext.java:331)
    at org.apache.hadoop.fs.FileContext.getFileContext(FileContext.java:448)
    at org.apache.spark.sql.delta.storage.HDFSLogStore.getFileContext(HDFSLogStore.scala:52)
    at org.apache.spark.sql.delta.storage.HDFSLogStore.read(HDFSLogStore.scala:56)
    at org.apache.spark.sql.delta.Checkpoints$class.loadMetadataFromFile(Checkpoints.scala:138)
    at org.apache.spark.sql.delta.Checkpoints$class.lastCheckpoint(Checkpoints.scala:132)
    at org.apache.spark.sql.delta.DeltaLog.lastCheckpoint(DeltaLog.scala:56)
    at org.apache.spark.sql.delta.DeltaLog.<init>(DeltaLog.scala:133)
    at org.apache.spark.sql.delta.DeltaLog$$anon$3$$anonfun$call$1$$anonfun$apply$7.apply(DeltaLog.scala:722)
    at org.apache.spark.sql.delta.DeltaLog$$anon$3$$anonfun$call$1$$anonfun$apply$7.apply(DeltaLog.scala:722)
    at org.apache.spark.sql.catalyst.plans.logical.AnalysisHelper$.allowInvokingTransformsInAnalyzer(AnalysisHelper.scala:194)
    at org.apache.spark.sql.delta.DeltaLog$$anon$3$$anonfun$call$1.apply(DeltaLog.scala:721)
    at org.apache.spark.sql.delta.DeltaLog$$anon$3$$anonfun$call$1.apply(DeltaLog.scala:721)
    at com.databricks.spark.util.DatabricksLogging$class.recordOperation(DatabricksLogging.scala:77)
    at org.apache.spark.sql.delta.DeltaLog$.recordOperation(DeltaLog.scala:623)
    at org.apache.spark.sql.delta.metering.DeltaLogging$class.recordDeltaOperation(DeltaLogging.scala:103)
    at org.apache.spark.sql.delta.DeltaLog$.recordDeltaOperation(DeltaLog.scala:623)
    at org.apache.spark.sql.delta.DeltaLog$$anon$3.call(DeltaLog.scala:720)
    at org.apache.spark.sql.delta.DeltaLog$$anon$3.call(DeltaLog.scala:718)
    at com.google.common.cache.LocalCache$LocalManualCache$1.load(LocalCache.java:4767)
    at com.google.common.cache.LocalCache$LoadingValueReference.loadFuture(LocalCache.java:3568)
    at com.google.common.cache.LocalCache$Segment.loadSync(LocalCache.java:2350)
    at com.google.common.cache.LocalCache$Segment.lockedGetOrLoad(LocalCache.java:2313)
    at com.google.common.cache.LocalCache$Segment.get(LocalCache.java:2228)
    at com.google.common.cache.LocalCache.get(LocalCache.java:3965)
    at com.google.common.cache.LocalCache$LocalManualCache.get(LocalCache.java:4764)
    at org.apache.spark.sql.delta.DeltaLog$.apply(DeltaLog.scala:718)
    at org.apache.spark.sql.delta.DeltaLog$.forTable(DeltaLog.scala:650)
    at org.apache.spark.sql.delta.sources.DeltaDataSource.createRelation(DeltaDataSource.scala:139)
    at org.apache.spark.sql.execution.datasources.SaveIntoDataSourceCommand.run(SaveIntoDataSourceCommand.scala:45)
    at org.apache.spark.sql.execution.command.ExecutedCommandExec.sideEffectResult$lzycompute(commands.scala:70)
    at org.apache.spark.sql.execution.command.ExecutedCommandExec.sideEffectResult(commands.scala:68)
    at org.apache.spark.sql.execution.command.ExecutedCommandExec.doExecute(commands.scala:86)
    at org.apache.spark.sql.execution.SparkPlan$$anonfun$execute$1.apply(SparkPlan.scala:131)
    at org.apache.spark.sql.execution.SparkPlan$$anonfun$execute$1.apply(SparkPlan.scala:127)
    at org.apache.spark.sql.execution.SparkPlan$$anonfun$executeQuery$1.apply(SparkPlan.scala:155)
    at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:151)
    at org.apache.spark.sql.execution.SparkPlan.executeQuery(SparkPlan.scala:152)
    at org.apache.spark.sql.execution.SparkPlan.execute(SparkPlan.scala:127)
    at org.apache.spark.sql.execution.QueryExecution.toRdd$lzycompute(QueryExecution.scala:80)
    at org.apache.spark.sql.execution.QueryExecution.toRdd(QueryExecution.scala:80)
    at org.apache.spark.sql.DataFrameWriter$$anonfun$runCommand$1.apply(DataFrameWriter.scala:676)
    at org.apache.spark.sql.DataFrameWriter$$anonfun$runCommand$1.apply(DataFrameWriter.scala:676)
    at org.apache.spark.sql.execution.SQLExecution$$anonfun$withNewExecutionId$1.apply(SQLExecution.scala:78)
    at org.apache.spark.sql.execution.SQLExecution$.withSQLConfPropagated(SQLExecution.scala:125)
    at org.apache.spark.sql.execution.SQLExecution$.withNewExecutionId(SQLExecution.scala:73)
    at org.apache.spark.sql.DataFrameWriter.runCommand(DataFrameWriter.scala:676)
    at org.apache.spark.sql.DataFrameWriter.saveToV1Source(DataFrameWriter.scala:285)
    at org.apache.spark.sql.DataFrameWriter.save(DataFrameWriter.scala:271)
    at org.apache.spark.sql.DataFrameWriter.save(DataFrameWriter.scala:229)
    at CatalogData$.main(CatalogData.scala:26)
    at CatalogData.main(CatalogData.scala)

I'm using the master build of delta-io. Is s3a file system not implemented in delta yet, or am I missing something?

Thanks!

Most helpful comment

Thanks for your help. This issue was coming because I was using the 2.6.0 version of hadoop-aws, which doesn't have AbstractFileSystem implementation for s3a. Upgrading it to 2.8.0 fixed the issue.

All 3 comments

You need to have the S3 file system implementation from the maven
dependency hadoop-aws
https://mvnrepository.com/artifact/org.apache.hadoop/hadoop-aws in the
classpath. Either add it in spark-submit or to your maven/sbt project file.

Note that the current release 0.1.0 does not correctly work with concurrent
writes to S3 because S3 filesystem implementation does not provide any way
to write-file-if-absent and so we cannot transactionally update a log. The
next release 0.2.0 (in a couple of days) is going to have support for
concurrent writes to S3 as long as all the writes go through a single Spark
driver application. I recommend using that when 0.2.0 is released. Along
with that, we will also have detailed docs on how to set up stuff for using
Delta on S3.

TD

On Tue, Jun 18, 2019 at 1:52 PM Gaurav Sehgal notifications@github.com
wrote:

Hi, I'm trying to use delta lake over s3 storage layer.

Spark command: df.write.format("delta").save("s3a:///")

Getting this error.

19/06/19 02:13:59 WARN DeltaLog: Failed to parse s3a://demo-atlan-source/delta_test/_delta_log/_last_checkpoint. This may happen if there was an error during read operation, or a file appears to be partial. Sleeping and trying again.
org.apache.hadoop.fs.UnsupportedFileSystemException: No AbstractFileSystem for scheme: s3a
at org.apache.hadoop.fs.AbstractFileSystem.createFileSystem(AbstractFileSystem.java:154)
at org.apache.hadoop.fs.AbstractFileSystem.get(AbstractFileSystem.java:242)
at org.apache.hadoop.fs.FileContext$2.run(FileContext.java:334)
at org.apache.hadoop.fs.FileContext$2.run(FileContext.java:331)
at java.security.AccessController.doPrivileged(Native Method)
at javax.security.auth.Subject.doAs(Subject.java:422)
at org.apache.hadoop.security.UserGroupInformation.doAs(UserGroupInformation.java:1692)
at org.apache.hadoop.fs.FileContext.getAbstractFileSystem(FileContext.java:331)
at org.apache.hadoop.fs.FileContext.getFileContext(FileContext.java:448)
at org.apache.spark.sql.delta.storage.HDFSLogStore.getFileContext(HDFSLogStore.scala:52)
at org.apache.spark.sql.delta.storage.HDFSLogStore.read(HDFSLogStore.scala:56)
at org.apache.spark.sql.delta.Checkpoints$class.loadMetadataFromFile(Checkpoints.scala:138)
at org.apache.spark.sql.delta.Checkpoints$class.lastCheckpoint(Checkpoints.scala:132)
at org.apache.spark.sql.delta.DeltaLog.lastCheckpoint(DeltaLog.scala:56)
at org.apache.spark.sql.delta.DeltaLog.(DeltaLog.scala:133)
at org.apache.spark.sql.delta.DeltaLog$$anon$3$$anonfun$call$1$$anonfun$apply$7.apply(DeltaLog.scala:722)
at org.apache.spark.sql.delta.DeltaLog$$anon$3$$anonfun$call$1$$anonfun$apply$7.apply(DeltaLog.scala:722)
at org.apache.spark.sql.catalyst.plans.logical.AnalysisHelper$.allowInvokingTransformsInAnalyzer(AnalysisHelper.scala:194)
at org.apache.spark.sql.delta.DeltaLog$$anon$3$$anonfun$call$1.apply(DeltaLog.scala:721)
at org.apache.spark.sql.delta.DeltaLog$$anon$3$$anonfun$call$1.apply(DeltaLog.scala:721)
at com.databricks.spark.util.DatabricksLogging$class.recordOperation(DatabricksLogging.scala:77)
at org.apache.spark.sql.delta.DeltaLog$.recordOperation(DeltaLog.scala:623)
at org.apache.spark.sql.delta.metering.DeltaLogging$class.recordDeltaOperation(DeltaLogging.scala:103)
at org.apache.spark.sql.delta.DeltaLog$.recordDeltaOperation(DeltaLog.scala:623)
at org.apache.spark.sql.delta.DeltaLog$$anon$3.call(DeltaLog.scala:720)
at org.apache.spark.sql.delta.DeltaLog$$anon$3.call(DeltaLog.scala:718)
at com.google.common.cache.LocalCache$LocalManualCache$1.load(LocalCache.java:4767)
at com.google.common.cache.LocalCache$LoadingValueReference.loadFuture(LocalCache.java:3568)
at com.google.common.cache.LocalCache$Segment.loadSync(LocalCache.java:2350)
at com.google.common.cache.LocalCache$Segment.lockedGetOrLoad(LocalCache.java:2313)
at com.google.common.cache.LocalCache$Segment.get(LocalCache.java:2228)
at com.google.common.cache.LocalCache.get(LocalCache.java:3965)
at com.google.common.cache.LocalCache$LocalManualCache.get(LocalCache.java:4764)
at org.apache.spark.sql.delta.DeltaLog$.apply(DeltaLog.scala:718)
at org.apache.spark.sql.delta.DeltaLog$.forTable(DeltaLog.scala:650)
at org.apache.spark.sql.delta.sources.DeltaDataSource.createRelation(DeltaDataSource.scala:139)
at org.apache.spark.sql.execution.datasources.SaveIntoDataSourceCommand.run(SaveIntoDataSourceCommand.scala:45)
at org.apache.spark.sql.execution.command.ExecutedCommandExec.sideEffectResult$lzycompute(commands.scala:70)
at org.apache.spark.sql.execution.command.ExecutedCommandExec.sideEffectResult(commands.scala:68)
at org.apache.spark.sql.execution.command.ExecutedCommandExec.doExecute(commands.scala:86)
at org.apache.spark.sql.execution.SparkPlan$$anonfun$execute$1.apply(SparkPlan.scala:131)
at org.apache.spark.sql.execution.SparkPlan$$anonfun$execute$1.apply(SparkPlan.scala:127)
at org.apache.spark.sql.execution.SparkPlan$$anonfun$executeQuery$1.apply(SparkPlan.scala:155)
at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:151)
at org.apache.spark.sql.execution.SparkPlan.executeQuery(SparkPlan.scala:152)
at org.apache.spark.sql.execution.SparkPlan.execute(SparkPlan.scala:127)
at org.apache.spark.sql.execution.QueryExecution.toRdd$lzycompute(QueryExecution.scala:80)
at org.apache.spark.sql.execution.QueryExecution.toRdd(QueryExecution.scala:80)
at org.apache.spark.sql.DataFrameWriter$$anonfun$runCommand$1.apply(DataFrameWriter.scala:676)
at org.apache.spark.sql.DataFrameWriter$$anonfun$runCommand$1.apply(DataFrameWriter.scala:676)
at org.apache.spark.sql.execution.SQLExecution$$anonfun$withNewExecutionId$1.apply(SQLExecution.scala:78)
at org.apache.spark.sql.execution.SQLExecution$.withSQLConfPropagated(SQLExecution.scala:125)
at org.apache.spark.sql.execution.SQLExecution$.withNewExecutionId(SQLExecution.scala:73)
at org.apache.spark.sql.DataFrameWriter.runCommand(DataFrameWriter.scala:676)
at org.apache.spark.sql.DataFrameWriter.saveToV1Source(DataFrameWriter.scala:285)
at org.apache.spark.sql.DataFrameWriter.save(DataFrameWriter.scala:271)
at org.apache.spark.sql.DataFrameWriter.save(DataFrameWriter.scala:229)
at CatalogData$.main(CatalogData.scala:26)
at CatalogData.main(CatalogData.scala)

I'm using the master build of delta-io. Is s3a file system not implemented
in delta yet, or am I missing something?

Thanks!

—
You are receiving this because you are subscribed to this thread.
Reply to this email directly, view it on GitHub
https://github.com/delta-io/delta/issues/68?email_source=notifications&email_token=AAFB5LBVPZTMVJ546PNUSADP3FDJPA5CNFSM4HZDYRSKYY3PNVWWK3TUL52HS4DFUVEXG43VMWVGG33NNVSW45C7NFSM4G2HZXDA,
or mute the thread
https://github.com/notifications/unsubscribe-auth/AAFB5LGDAMAXUFTNAVL7AVTP3FDJPANCNFSM4HZDYRSA
.

Thanks for your help. This issue was coming because I was using the 2.6.0 version of hadoop-aws, which doesn't have AbstractFileSystem implementation for s3a. Upgrading it to 2.8.0 fixed the issue.

Awesome! Glad it got sorted out.

On Tue, Jun 18, 2019 at 2:15 PM Gaurav Sehgal notifications@github.com
wrote:

Thanks for your help. This issue was coming because I was using the 2.6.0
version of hadoop-aws, which doesn't have AbstractFileSystem
implementation for s3a. Upgrading it to 2.8.0 fixed the issue.

—
You are receiving this because you commented.
Reply to this email directly, view it on GitHub
https://github.com/delta-io/delta/issues/68?email_source=notifications&email_token=AAFB5LAINROMYUHLWIYSUBLP3FF7NA5CNFSM4HZDYRSKYY3PNVWWK3TUL52HS4DFVREXG43VMVBW63LNMVXHJKTDN5WW2ZLOORPWSZGODX77UBY#issuecomment-503314951,
or mute the thread
https://github.com/notifications/unsubscribe-auth/AAFB5LAMPNVQARQ7JUFUIM3P3FF7NANCNFSM4HZDYRSA
.

Was this page helpful?
0 / 5 - 0 ratings

Related issues

NimeshSatam picture NimeshSatam  Â·  3Comments

SriramAvatar picture SriramAvatar  Â·  4Comments

marmbrus picture marmbrus  Â·  7Comments

NimeshSatam picture NimeshSatam  Â·  8Comments

pyMixin picture pyMixin  Â·  5Comments