Skip to content

[SPARK-19926][PYSPARK] Make pyspark exception more user-friendly - #17267

Closed
uncleGen wants to merge 3 commits into
apache:masterfrom
uncleGen:SPARK-19926
Closed

[SPARK-19926][PYSPARK] Make pyspark exception more user-friendly#17267
uncleGen wants to merge 3 commits into
apache:masterfrom
uncleGen:SPARK-19926

Conversation

@uncleGen

@uncleGenuncleGen commented Mar 12, 2017

Copy link
Copy Markdown
Contributor

What changes were proposed in this pull request?

Exception in pyspark is a little difficult to read.

before pr, like:

Traceback (most recent call last):
File "<stdin>", line 5, in <module>
File "/root/dev/spark/dist/python/pyspark/sql/streaming.py", line 853, in start
return self._sq(self._jwrite.start())
File "/root/dev/spark/dist/python/lib/py4j-0.10.4-src.zip/py4j/java_gateway.py", line 1133, in __call__
File "/root/dev/spark/dist/python/pyspark/sql/utils.py", line 69, in deco
raise AnalysisException(s.split(': ', 1)[1], stackTrace)
pyspark.sql.utils.AnalysisException: u'Append output mode not supported when there are streaming aggregations on streaming DataFrames/DataSets without watermark;;\nAggregate [window#17, word#5], [window#17 AS window#11, word#5, count(1) AS count#16L]\n+- Filter ((t#6 >= window#17.start) && (t#6 < window#17.end))\n +- Expand [ArrayBuffer(named_struct(start, ((((CEIL((cast((precisetimestamp(t#6) - 0) as double) / cast(30000000 as double))) + cast(0 as bigint)) - cast(1 as bigint)) * 30000000) + 0), end, (((((CEIL((cast((precisetimestamp(t#6) - 0) as double) / cast(30000000 as double))) + cast(0 as bigint)) - cast(1 as bigint)) * 30000000) + 0) + 30000000)), word#5, t#6-T30000ms), ArrayBuffer(named_struct(start, ((((CEIL((cast((precisetimestamp(t#6) - 0) as double) / cast(30000000 as double))) + cast(1 as bigint)) - cast(1 as bigint)) * 30000000) + 0), end, (((((CEIL((cast((precisetimestamp(t#6) - 0) as double) / cast(30000000 as double))) + cast(1 as bigint)) - cast(1 as bigint)) * 30000000) + 0) + 30000000)), word#5, t#6-T30000ms)], [window#17, word#5, t#6-T30000ms]\n +- EventTimeWatermark t#6: timestamp, interval 30 seconds\n +- Project [cast(word#0 as string) AS word#5, cast(t#1 as timestamp) AS t#6]\n +- StreamingRelation DataSource(org.apache.spark.sql.SparkSession@c4079ca,csv,List(),Some(StructType(StructField(word,StringType,true), StructField(t,IntegerType,true))),List(),None,Map(sep -> ;, path -> /tmp/data),None), FileSource[/tmp/data], [word#0, t#1]\n'

after pr:

Traceback (most recent call last):
File "<stdin>", line 5, in <module>
File "/root/dev/spark/dist/python/pyspark/sql/streaming.py", line 853, in start
return self._sq(self._jwrite.start())
File "/root/dev/spark/dist/python/lib/py4j-0.10.4-src.zip/py4j/java_gateway.py", line 1133, in __call__
File "/root/dev/spark/dist/python/pyspark/sql/utils.py", line 69, in deco
raise AnalysisException(s.split(': ', 1)[1], stackTrace)
pyspark.sql.utils.AnalysisException: Append output mode not supported when there are streaming aggregations on streaming DataFrames/DataSets without watermark;;
Aggregate [window#17, word#5], [window#17 AS window#11, word#5, count(1) AS count#16L]
+- Filter ((t#6 >= window#17.start) && (t#6 < window#17.end))
+- Expand [ArrayBuffer(named_struct(start, ((((CEIL((cast((precisetimestamp(t#6) - 0) as double) / cast(30000000 as double))) + cast(0 as bigint)) - cast(1 as bigint)) * 30000000) + 0), end, (((((CEIL((cast((precisetimestamp(t#6) - 0) as double) / cast(30000000 as double))) + cast(0 as bigint)) - cast(1 as bigint)) * 30000000) + 0) + 30000000)), word#5, t#6-T30000ms), ArrayBuffer(named_struct(start, ((((CEIL((cast((precisetimestamp(t#6) - 0) as double) / cast(30000000 as double))) + cast(1 as bigint)) - cast(1 as bigint)) * 30000000) + 0), end, (((((CEIL((cast((precisetimestamp(t#6) - 0) as double) / cast(30000000 as double))) + cast(1 as bigint)) - cast(1 as bigint)) * 30000000) + 0) + 30000000)), word#5, t#6-T30000ms)], [window#17, word#5, t#6-T30000ms]
+- EventTimeWatermark t#6: timestamp, interval 30 seconds
+- Project [cast(word#0 as string) AS word#5, cast(t#1 as timestamp) AS t#6]
+- StreamingRelation DataSource(org.apache.spark.sql.SparkSession@5265083b,csv,List(),Some(StructType(StructField(word,StringType,true), StructField(t,IntegerType,true))),List(),None,Map(sep -> ;, path -> /tmp/data),None), FileSource[/tmp/data], [word#0, t#1]

IMHO, the root cause is the repr is not user-friendly

class CapturedException(Exception):
def __init__(self, desc, stackTrace):
self.desc = desc
self.stackTrace = stackTrace
def __str__(self):
return repr(self.desc)
▲▲▲▲

This pr change repr to str

str()repr()
make object readableneed code that reproduces object
generate output for end usergenerate output for developer

How was this patch tested?

Jenkins

@uncleGen

uncleGen commented Mar 12, 2017

Copy link
Copy Markdown
ContributorAuthor

Maybe @viirya and @davies can give some suggestion.

@SparkQA

Copy link
Copy Markdown

Test build #74403 has finished for PR 17267 at commit 273c1bc.

  • This patch passes all tests.
  • This patch merges cleanly.
  • This patch adds no public classes.

@viirya

Copy link
Copy Markdown
Member

Thanks for working on this. LGTM.

@uncleGen

uncleGen commented Mar 13, 2017

Copy link
Copy Markdown
ContributorAuthor

@viirya Thanks for your review.
cc @srowen IIUC, it is OK for other use of repr.

@srowen

Copy link
Copy Markdown
Member

What's the difference between the two? briefly. I don't know enough to evaluate it though the effect looks positive. Is this the only place this should change?

@uncleGen

uncleGen commented Mar 13, 2017

Copy link
Copy Markdown
ContributorAuthor

IMHO, yes. And @viirya is the original author.

str()repr()
make object readableneed code that reproduces object
generate output for end usergenerate output for developer

Comment threadpython/pyspark/sql/utils.py Outdated

def __str__(self):
return repr(self.desc)
return str(self.desc)

@HyukjinKwonHyukjinKwonMar 13, 2017

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Hm.. does this work for unicode in Python 2, for example, spark.range(1).select("아")? Up to my knowledge, converting it to ascii directly throws an exception.

>>>str(u"아")
Traceback (most recent call last):
File "<stdin>", line 1, in <module>
UnicodeEncodeError: 'ascii' codec can't encode character u'\uc544' in position 0: ordinal not in range(128)
>>>repr(u"아")
"u'\\uc544'"

Maybe, we should check if this is unicode and do .encode.

@HyukjinKwonHyukjinKwonMar 13, 2017

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I just tested with this change as below to help:

  • before
>>>try:
... spark.range(1).select(u"아")
... exceptExceptionase:
... printe
u"cannot resolve '`\uc544`' given input columns: [id];;\n'Project ['\uc544]\n+- Range (0, 1, step=1, splits=Some(8))\n"
>>>spark.range(1).select(u"아")
Traceback (most recent call last):
File "<stdin>", line 1, in <module>
File ".../spark/python/pyspark/sql/dataframe.py", line 992, in select
jdf = self._jdf.select(self._jcols(*cols))
File ".../spark/python/lib/py4j-0.10.4-src.zip/py4j/java_gateway.py", line 1133, in __call__
File ".../spark/python/pyspark/sql/utils.py", line 69, in deco
raise AnalysisException(s.split(': ', 1)[1], stackTrace)
pyspark.sql.utils.AnalysisException: u"cannot resolve '`\uc544`' given input columns: [id];;\n'Project ['\uc544]\n+- Range (0, 1, step=1, splits=Some(8))\n"
  • after
>>>try:
... spark.range(1).select(u"아")
... exceptExceptionase:
... printe
Traceback (most recent call last):
File "<stdin>", line 4, in <module>
File ".../spark/python/pyspark/sql/utils.py", line 27, in __str__
return str(self.desc)
UnicodeEncodeError: 'ascii' codec can't encode character u'\uc544' in position 17: ordinal not in range(128)
>>>spark.range(1).select(u"아")
Traceback (most recent call last):
File "<stdin>", line 1, in <module>
File ".../spark/python/pyspark/sql/dataframe.py", line 992, in select
jdf = self._jdf.select(self._jcols(*cols))
File ".../spark/python/lib/py4j-0.10.4-src.zip/py4j/java_gateway.py", line 1133, in __call__
File ".../spark/python/pyspark/sql/utils.py", line 69, in deco
raise AnalysisException(s.split(': ', 1)[1], stackTrace)
pyspark.sql.utils.AnalysisException

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@uncleGen, could you double check if I did something wrong maybe?

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

We can add a check under Python2. If it is unicode, just encode it with utf-8.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@HyukjinKwon Good catch!

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Ah, thank you for confirmation. I thought I was mistaken :).

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Maybe another benefit for this change is, before it you will see the error log in your example like:

u"cannot resolve '\uc544' given input columns: [id];;\n'Project ['\uc544]

repr will show unicode escape characters \uc544. Even you encode it, you will see binary representation for it. str can show the correct "아" if encoded with utf-8.

If I test it correctly.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yea, I support this change and tested some more cases with that encode.

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

based on latest commit:

>>> df.select("아")
Traceback (most recent call last):
File "<stdin>", line 1, in <module>
File ".../spark/python/pyspark/sql/dataframe.py", line 992, in select
jdf = self._jdf.select(self._jcols(*cols))
File ".../spark/python/lib/py4j-0.10.4-src.zip/py4j/java_gateway.py", line 1133, in __call__
File ".../spark/python/pyspark/sql/utils.py", line 75, in deco
raise AnalysisException(s.split(': ', 1)[1], stackTrace)
pyspark.sql.utils.AnalysisException
: cannot resolve '`아`' given input columns: [age, name];;
'Project ['아]
+- Relation[age#0L,name#1] json

@uncleGen

Copy link
Copy Markdown
ContributorAuthor

Thanks @HyukjinKwon,you give a good catch!I lost that case. Thanks @viirya for your suggestion.

@SparkQA

Copy link
Copy Markdown

Test build #74487 has finished for PR 17267 at commit 5bc1d8e.

  • This patch fails PySpark unit tests.
  • This patch merges cleanly.
  • This patch adds no public classes.

@SparkQA

Copy link
Copy Markdown

Test build #74490 has finished for PR 17267 at commit 6c55e02.

  • This patch fails Python style tests.
  • This patch merges cleanly.
  • This patch adds no public classes.

@SparkQA

Copy link
Copy Markdown

Test build #74491 has finished for PR 17267 at commit 7b96e97.

  • This patch passes all tests.
  • This patch merges cleanly.
  • This patch adds no public classes.

Comment threadpython/pyspark/sql/utils.py Outdated
import py4j
import sys

if sys.version > '3':

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I think it should be >=.

@SparkQA

Copy link
Copy Markdown

Test build #74503 has finished for PR 17267 at commit edf9b12.

  • This patch passes all tests.
  • This patch merges cleanly.
  • This patch adds no public classes.

@uncleGen

uncleGen commented Mar 15, 2017

Copy link
Copy Markdown
ContributorAuthor

ping @viirya and @HyukjinKwon

@viirya

Copy link
Copy Markdown
Member

LGTM cc @davies@holdenk

@uncleGen

Copy link
Copy Markdown
ContributorAuthor

@srowen Could you please take a view and help to merge?

@srowen

Copy link
Copy Markdown
Member

I'm not reviewing this patch. People who know better should merge it

@holdenk

Copy link
Copy Markdown
Contributor

I'll take a look at reviewing this later on this week @uncleGen. Two minor thing that we can do in the meantime is make the JIRA description a bit clearer as to what the proposed change is, the other is this change isn't really tested by Jenkins - there are no tests that look at the formatting of the error strings - maybe consider adding a test or updating the description on the PR.

@uncleGenuncleGen changed the title [SPARK-19926][PYSPARK] Make pyspark exception more readable[SPARK-19926][PYSPARK] Make pyspark exception more user-friendlyMar 21, 2017
@gatorsmile

Copy link
Copy Markdown
Member

cc @ueshin

@HyukjinKwon

Copy link
Copy Markdown
Member

LGTM too but hope there would be a test if possible.

@ueshin

Copy link
Copy Markdown
Member

Correct me if I'm wrong, but I got the following message after this patch in Python 3.6:

Traceback (mostrecentcalllast):
File"<stdin>", line1, in<module>File"/Users/ueshin/workspace/pyspark/spark/python/pyspark/sql/dataframe.py", line1049, inselectjdf=self._jdf.select(self._jcols(*cols))
File"/Users/ueshin/workspace/pyspark/spark/python/lib/py4j-0.10.4-src.zip/py4j/java_gateway.py", line1133, in__call__File"/Users/ueshin/workspace/pyspark/spark/python/pyspark/sql/utils.py", line77, indecoraiseAnalysisException(s.split(': ', 1)[1], stackTrace)
pyspark.sql.utils.AnalysisException: b"cannot resolve '`\xec\x95\x84`' given input columns: [id];;\n'Project ['\xec\x95\x84]\n+- Range (0, 1, step=1, splits=Some(8))\n"

I guess this message is not desirable?

@ueshin

Copy link
Copy Markdown
Member

+1 for adding a test.

return repr(self.desc)
desc = self.desc
if isinstance(desc, unicode):
return str(desc.encode('utf-8'))

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@ueshin, you are right and I misread the codes. We need to

  • unicode in Python 2 => u.encode("utf-8").
  • others in Python 2 => return str(s).
  • others in Python 3 => return str(s).

Root cause for #17267 (comment) looks because encode on string (also same as unicode in Python 2) in Python 3 produces 8-bit bytes, b"...", (also same as normal string, "..." and b"...", where b is ignored, in Python 2). And str function works differently as below:

Python 2

>>>str(b"aa")
'aa'>>>b"aa"'aa'

Python 3

>>>str(b"aa")
"b'aa'">>>"aa"'aa'

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Good catch! I previously thought str works like Python2.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

cc @zero323 and @davies too. Would you have some time to take a look for this one? This is a typical annoying problem between unicode and byte strings. There are many similar PRs (at least I can identify few PRs trying to handle this problem. One good example might help resolving other PRs too.

@viirya

Copy link
Copy Markdown
Member

+1 We should add a test for this.

@holdenk

Copy link
Copy Markdown
Contributor

Hey @uncleGen anytime to add a test for this?

@HyukjinKwonHyukjinKwon mentioned this pull request Jul 31, 2017
@asfgitasfgit closed this in 3a45c7fAug 5, 2017
@jiangxb1987

Copy link
Copy Markdown
Contributor

@dataknocker do you want to take over this one? then we can continue with #18324

HyukjinKwon pushed a commit that referenced this pull request Sep 18, 2019
…endly
### What changes were proposed in this pull request?
The str of `CapaturedException` is now returned by str(self.desc) rather than repr(self.desc), which is more user-friendly. It also handles unicode under python2 specially.
### Why are the changes needed?
This is an improvement, and makes exception more human readable in python side.
### Does this PR introduce any user-facing change?
Before this pr, select `中文字段` throws exception something likes below:
```
Traceback (most recent call last):
File "/Users/advancedxy/code_workspace/github/spark/python/pyspark/sql/tests/test_utils.py", line 34, in test_capture_user_friendly_exception
raise e
AnalysisException: u"cannot resolve '`\u4e2d\u6587\u5b57\u6bb5`' given input columns: []; line 1 pos 7;\n'Project ['\u4e2d\u6587\u5b57\u6bb5]\n+- OneRowRelation\n"
```
after this pr:
```
Traceback (most recent call last):
File "/Users/advancedxy/code_workspace/github/spark/python/pyspark/sql/tests/test_utils.py", line 34, in test_capture_user_friendly_exception
raise e
AnalysisException: cannot resolve '`中文字段`' given input columns: []; line 1 pos 7;
'Project ['中文字段]
+- OneRowRelation
```
### How was this patch
Add a new test to verify unicode are correctly converted and manual checks for thrown exceptions.
This pr's credits should go to uncleGen and is based on #17267Closes#25814 from advancedxy/python_exception_19926_and_21045.
Authored-by: Xianjin YE <advancedxy@gmail.com>
Signed-off-by: HyukjinKwon <gurwls223@apache.org>
zifeif2 pushed a commit to zifeif2/spark that referenced this pull request Nov 22, 2025
## What changes were proposed in this pull request?
This PR proposes to close stale PRs, mostly the same instances with apache#18017Closesapache#14085 - [SPARK-16408][SQL] SparkSQL Added file get Exception: is a directory …
Closesapache#14239 - [SPARK-16593] [CORE] [WIP] Provide a pre-fetch mechanism to accelerate shuffle stage.
Closesapache#14567 - [SPARK-16992][PYSPARK] Python Pep8 formatting and import reorganisation
Closesapache#14579 - [SPARK-16921][PYSPARK] RDD/DataFrame persist()/cache() should return Python context managers
Closesapache#14601 - [SPARK-13979][Core] Killed executor is re spawned without AWS key…
Closesapache#14830 - [SPARK-16992][PYSPARK][DOCS] import sort and autopep8 on Pyspark examples
Closesapache#14963 - [SPARK-16992][PYSPARK] Virtualenv for Pylint and pep8 in lint-python
Closesapache#15227 - [SPARK-17655][SQL]Remove unused variables declarations and definations in a WholeStageCodeGened stage
Closesapache#15240 - [SPARK-17556] [CORE] [SQL] Executor side broadcast for broadcast joins
Closesapache#15405 - [SPARK-15917][CORE] Added support for number of executors in Standalone [WIP]
Closesapache#16099 - [SPARK-18665][SQL] set statement state to "ERROR" after user cancel job
Closesapache#16445 - [SPARK-19043][SQL]Make SparkSQLSessionManager more configurable
Closesapache#16618 - [SPARK-14409][ML][WIP] Add RankingEvaluator
Closesapache#16766 - [SPARK-19426][SQL] Custom coalesce for Dataset
Closesapache#16832 - [SPARK-19490][SQL] ignore case sensitivity when filtering hive partition columns
Closesapache#17052 - [SPARK-19690][SS] Join a streaming DataFrame with a batch DataFrame which has an aggregation may not work
Closesapache#17267 - [SPARK-19926][PYSPARK] Make pyspark exception more user-friendly
Closesapache#17371 - [SPARK-19903][PYSPARK][SS] window operator miss the `watermark` metadata of time column
Closesapache#17401 - [SPARK-18364][YARN] Expose metrics for YarnShuffleService
Closesapache#17519 - [SPARK-15352][Doc] follow-up: add configuration docs for topology-aware block replication
Closesapache#17530 - [SPARK-5158] Access kerberized HDFS from Spark standalone
Closesapache#17854 - [SPARK-20564][Deploy] Reduce massive executor failures when executor count is large (>2000)
Closesapache#17979 - [SPARK-19320][MESOS][WIP]allow specifying a hard limit on number of gpus required in each spark executor when running on mesos
Closesapache#18127 - [SPARK-6628][SQL][Branch-2.1] Fix ClassCastException when executing sql statement 'insert into' on hbase table
Closesapache#18236 - [SPARK-21015] Check field name is not null and empty in GenericRowWit…
Closesapache#18269 - [SPARK-21056][SQL] Use at most one spark job to list files in InMemoryFileIndex
Closesapache#18328 - [SPARK-21121][SQL] Support changing storage level via the spark.sql.inMemoryColumnarStorage.level variable
Closesapache#18354 - [SPARK-18016][SQL][CATALYST][BRANCH-2.1] Code Generation: Constant Pool Limit - Class Splitting
Closesapache#18383 - [SPARK-21167][SS] Set kafka clientId while fetch messages
Closesapache#18414 - [SPARK-21169] [core] Make sure to update application status to RUNNING if executors are accepted and RUNNING after recovery
Closesapache#18432 - resolve com.esotericsoftware.kryo.KryoException
Closesapache#18490 - [SPARK-21269][Core][WIP] Fix FetchFailedException when enable maxReqSizeShuffleToMem and KryoSerializer
Closesapache#18585 - SPARK-21359
Closesapache#18609 - Spark SQL merge small files to big files Update InsertIntoHiveTable.scala
Added:
Closesapache#18308 - [SPARK-21099][Spark Core] INFO Log Message Using Incorrect Executor I…
Closesapache#18599 - [SPARK-21372] spark writes one log file even I set the number of spark_rotate_log to 0
Closesapache#18619 - [SPARK-21397][BUILD]Maven shade plugin adding dependency-reduced-pom.xml to …
Closesapache#18667 - Fix the simpleString used in error messages
Closesapache#18782 - Branch 2.1
Added:
Closesapache#17694 - [SPARK-12717][PYSPARK] Resolving race condition with pyspark broadcasts when using multiple threads
Added:
Closesapache#16456 - [SPARK-18994] clean up the local directories for application in future by annother thread
Closesapache#18683 - [SPARK-21474][CORE] Make number of parallel fetches from a reducer configurable
Closesapache#18690 - [SPARK-21334][CORE] Add metrics reporting service to External Shuffle Server
Added:
Closesapache#18827 - Merge pull request 1 from apache/master
## How was this patch tested?
N/A
Author: hyukjinkwon <gurwls223@gmail.com>
Closesapache#18780 from HyukjinKwon/close-prs.
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

9 participants

@uncleGen@SparkQA@viirya@srowen@holdenk@gatorsmile@HyukjinKwon@ueshin@jiangxb1987