Skip to content
Snippets Groups Projects
  1. Jul 19, 2017
    • jinxing's avatar
      [SPARK-21414] Refine SlidingWindowFunctionFrame to avoid OOM. · 4eb081cc
      jinxing authored
      ## What changes were proposed in this pull request?
      
      In `SlidingWindowFunctionFrame`, it is now adding all rows to the buffer for which the input row value is equal to or less than the output row upper bound, then drop all rows from the buffer for which the input row value is smaller than the output row lower bound.
      This could result in the buffer is very big though the window is small.
      For example:
      ```
      select a, b, sum(a)
      over (partition by b order by a range between 1000000 following and 1000001 following)
      from table
      ```
      We can refine the logic and just add the qualified rows into buffer.
      
      ## How was this patch tested?
      Manual test:
      Run sql
      `select shop, shopInfo, district, sum(revenue) over(partition by district order by revenue range between 100 following and 200 following) from revenueList limit 10`
      against a table with 4  columns(shop: String, shopInfo: String, district: String, revenue: Int). The biggest partition is around 2G bytes, containing 200k lines.
      Configure the executor with 2G bytes memory.
      With the change in this pr, it works find. Without this change, below exception will be thrown.
      ```
      MemoryError: Java heap space
      	at org.apache.spark.sql.catalyst.expressions.UnsafeRow.copy(UnsafeRow.java:504)
      	at org.apache.spark.sql.catalyst.expressions.UnsafeRow.copy(UnsafeRow.java:62)
      	at org.apache.spark.sql.execution.window.SlidingWindowFunctionFrame.write(WindowFunctionFrame.scala:201)
      	at org.apache.spark.sql.execution.window.WindowExec$$anonfun$14$$anon$1.next(WindowExec.scala:365)
      	at org.apache.spark.sql.execution.window.WindowExec$$anonfun$14$$anon$1.next(WindowExec.scala:289)
      	at org.apache.spark.sql.catalyst.expressions.GeneratedClass$GeneratedIterator.processNext(Unknown Source)
      	at org.apache.spark.sql.execution.BufferedRowIterator.hasNext(BufferedRowIterator.java:43)
      	at org.apache.spark.sql.execution.WholeStageCodegenExec$$anonfun$8$$anon$1.hasNext(WholeStageCodegenExec.scala:395)
      	at org.apache.spark.sql.execution.SparkPlan$$anonfun$2.apply(SparkPlan.scala:231)
      	at org.apache.spark.sql.execution.SparkPlan$$anonfun$2.apply(SparkPlan.scala:225)
      	at org.apache.spark.rdd.RDD$$anonfun$mapPartitionsInternal$1$$anonfun$apply$25.apply(RDD.scala:827)
      	at org.apache.spark.rdd.RDD$$anonfun$mapPartitionsInternal$1$$anonfun$apply$25.apply(RDD.scala:827)
      	at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:38)
      	at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:323)
      	at org.apache.spark.rdd.RDD.iterator(RDD.scala:287)
      	at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:87)
      	at org.apache.spark.scheduler.Task.run(Task.scala:108)
      	at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:341)
      	at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1142)
      	at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:617)
      ```
      
      Author: jinxing <jinxing6042@126.com>
      
      Closes #18634 from jinxing64/SPARK-21414.
      4eb081cc
  2. Jul 18, 2017
    • xuanyuanking's avatar
      [SPARK-21435][SQL] Empty files should be skipped while write to file · 81c99a5b
      xuanyuanking authored
      ## What changes were proposed in this pull request?
      
      Add EmptyDirectoryWriteTask for empty task while writing files. Fix the empty result for parquet format by leaving the first partition for meta writing.
      
      ## How was this patch tested?
      
      Add new test in `FileFormatWriterSuite `
      
      Author: xuanyuanking <xyliyuanjian@gmail.com>
      
      Closes #18654 from xuanyuanking/SPARK-21435.
      81c99a5b
    • Tathagata Das's avatar
      [SPARK-21462][SS] Added batchId to StreamingQueryProgress.json · 84f1b25f
      Tathagata Das authored
      ## What changes were proposed in this pull request?
      
      - Added batchId to StreamingQueryProgress.json as that was missing from the generated json.
      - Also, removed recently added numPartitions from StatefulOperatorProgress as this value does not change through the query run, and there are other ways to find that.
      
      ## How was this patch tested?
      Updated unit tests
      
      Author: Tathagata Das <tathagata.das1565@gmail.com>
      
      Closes #18675 from tdas/SPARK-21462.
      84f1b25f
    • Sean Owen's avatar
      [SPARK-21415] Triage scapegoat warnings, part 1 · e26dac5f
      Sean Owen authored
      ## What changes were proposed in this pull request?
      
      Address scapegoat warnings for:
      - BigDecimal double constructor
      - Catching NPE
      - Finalizer without super
      - List.size is O(n)
      - Prefer Seq.empty
      - Prefer Set.empty
      - reverse.map instead of reverseMap
      - Type shadowing
      - Unnecessary if condition.
      - Use .log1p
      - Var could be val
      
      In some instances like Seq.empty, I avoided making the change even where valid in test code to keep the scope of the change smaller. Those issues are concerned with performance and it won't matter for tests.
      
      ## How was this patch tested?
      
      Existing tests
      
      Author: Sean Owen <sowen@cloudera.com>
      
      Closes #18635 from srowen/Scapegoat1.
      e26dac5f
  3. Jul 17, 2017
    • Tathagata Das's avatar
      [SPARK-21409][SS] Follow up PR to allow different types of custom metrics to be exposed · e9faae13
      Tathagata Das authored
      ## What changes were proposed in this pull request?
      
      Implementation may expose both timing as well as size metrics. This PR enables that.
      
      Author: Tathagata Das <tathagata.das1565@gmail.com>
      
      Closes #18661 from tdas/SPARK-21409-2.
      e9faae13
    • Tathagata Das's avatar
      [SPARK-21409][SS] Expose state store memory usage in SQL metrics and progress updates · 9d8c8317
      Tathagata Das authored
      ## What changes were proposed in this pull request?
      
      Currently, there is no tracking of memory usage of state stores. This JIRA is to expose that through SQL metrics and StreamingQueryProgress.
      
      Additionally, added the ability to expose implementation-specific metrics through the StateStore APIs to the SQLMetrics.
      
      ## How was this patch tested?
      Added unit tests.
      
      Author: Tathagata Das <tathagata.das1565@gmail.com>
      
      Closes #18629 from tdas/SPARK-21409.
      9d8c8317
    • gatorsmile's avatar
      [SPARK-21354][SQL] INPUT FILE related functions do not support more than one sources · e398c281
      gatorsmile authored
      ### What changes were proposed in this pull request?
      The build-in functions `input_file_name`, `input_file_block_start`, `input_file_block_length` do not support more than one sources, like what Hive does. Currently, Spark does not block it and the outputs are ambiguous/non-deterministic. It could be from any side.
      
      ```
      hive> select *, INPUT__FILE__NAME FROM t1, t2;
      FAILED: SemanticException Column INPUT__FILE__NAME Found in more than One Tables/Subqueries
      ```
      
      This PR blocks it and issues an error.
      
      ### How was this patch tested?
      Added a test case
      
      Author: gatorsmile <gatorsmile@gmail.com>
      
      Closes #18580 from gatorsmile/inputFileName.
      e398c281
  4. Jul 16, 2017
  5. Jul 14, 2017
  6. Jul 13, 2017
    • Sean Owen's avatar
      [SPARK-19810][BUILD][CORE] Remove support for Scala 2.10 · 425c4ada
      Sean Owen authored
      ## What changes were proposed in this pull request?
      
      - Remove Scala 2.10 build profiles and support
      - Replace some 2.10 support in scripts with commented placeholders for 2.12 later
      - Remove deprecated API calls from 2.10 support
      - Remove usages of deprecated context bounds where possible
      - Remove Scala 2.10 workarounds like ScalaReflectionLock
      - Other minor Scala warning fixes
      
      ## How was this patch tested?
      
      Existing tests
      
      Author: Sean Owen <sowen@cloudera.com>
      
      Closes #17150 from srowen/SPARK-19810.
      425c4ada
  7. Jul 12, 2017
    • Wenchen Fan's avatar
      [SPARK-17701][SQL] Refactor RowDataSourceScanExec so its sameResult call does not compare strings · 780586a9
      Wenchen Fan authored
      ## What changes were proposed in this pull request?
      
      Currently, `RowDataSourceScanExec` and `FileSourceScanExec` rely on a "metadata" string map to implement equality comparison, since the RDDs they depend on cannot be directly compared. This has resulted in a number of correctness bugs around exchange reuse, e.g. SPARK-17673 and SPARK-16818.
      
      To make these comparisons less brittle, we should refactor these classes to compare constructor parameters directly instead of relying on the metadata map.
      
      This PR refactors `RowDataSourceScanExec`, `FileSourceScanExec` will be fixed in the follow-up PR.
      
      ## How was this patch tested?
      
      existing tests
      
      Author: Wenchen Fan <wenchen@databricks.com>
      
      Closes #18600 from cloud-fan/minor.
      780586a9
    • liuxian's avatar
      [SPARK-21007][SQL] Add SQL function - RIGHT && LEFT · aaad34dc
      liuxian authored
      ## What changes were proposed in this pull request?
       Add  SQL function - RIGHT && LEFT, same as MySQL:
      https://dev.mysql.com/doc/refman/5.7/en/string-functions.html#function_left
      https://dev.mysql.com/doc/refman/5.7/en/string-functions.html#function_right
      
      ## How was this patch tested?
      unit test
      
      Author: liuxian <liu.xian3@zte.com.cn>
      
      Closes #18228 from 10110346/lx-wip-0607.
      aaad34dc
    • Burak Yavuz's avatar
      [SPARK-21370][SS] Add test for state reliability when one read-only state... · e0af76a3
      Burak Yavuz authored
      [SPARK-21370][SS] Add test for state reliability when one read-only state store aborts after read-write state store commits
      
      ## What changes were proposed in this pull request?
      
      During Streaming Aggregation, we have two StateStores per task, one used as read-only in
      `StateStoreRestoreExec`, and one read-write used in `StateStoreSaveExec`. `StateStore.abort`
      will be called for these StateStores if they haven't committed their results. We need to
      make sure that `abort` in read-only store after a `commit` in the read-write store doesn't
      accidentally lead to the deletion of state.
      
      This PR adds a test for this condition.
      
      ## How was this patch tested?
      
      This PR adds a test.
      
      Author: Burak Yavuz <brkyvz@gmail.com>
      
      Closes #18603 from brkyvz/ss-test.
      e0af76a3
    • Jane Wang's avatar
      [SPARK-12139][SQL] REGEX Column Specification · 2cbfc975
      Jane Wang authored
      ## What changes were proposed in this pull request?
      Hive interprets regular expression, e.g., `(a)?+.+` in query specification. This PR enables spark to support this feature when hive.support.quoted.identifiers is set to true.
      
      ## How was this patch tested?
      
      - Add unittests in SQLQuerySuite.scala
      - Run spark-shell tested the original failed query:
      scala> hc.sql("SELECT `(a|b)?+.+` from test1").collect.foreach(println)
      
      Author: Jane Wang <janewang@fb.com>
      
      Closes #18023 from janewangfb/support_select_regex.
      2cbfc975
  8. Jul 11, 2017
    • gatorsmile's avatar
      [SPARK-19285][SQL] Implement UDF0 · d3e07165
      gatorsmile authored
      ### What changes were proposed in this pull request?
      This PR is to implement UDF0. `UDF0` is needed when users need to implement a JAVA UDF with no argument.
      
      ### How was this patch tested?
      Added a test case
      
      Author: gatorsmile <gatorsmile@gmail.com>
      
      Closes #18598 from gatorsmile/udf0.
      d3e07165
    • hyukjinkwon's avatar
      [SPARK-21365][PYTHON] Deduplicate logics parsing DDL type/schema definition · ebc124d4
      hyukjinkwon authored
      ## What changes were proposed in this pull request?
      
      This PR deals with four points as below:
      
      - Reuse existing DDL parser APIs rather than reimplementing within PySpark
      
      - Support DDL formatted string, `field type, field type`.
      
      - Support case-insensitivity for parsing.
      
      - Support nested data types as below:
      
        **Before**
        ```
        >>> spark.createDataFrame([[[1]]], "struct<a: struct<b: int>>").show()
        ...
        ValueError: The strcut field string format is: 'field_name:field_type', but got: a: struct<b: int>
        ```
      
        ```
        >>> spark.createDataFrame([[[1]]], "a: struct<b: int>").show()
        ...
        ValueError: The strcut field string format is: 'field_name:field_type', but got: a: struct<b: int>
        ```
      
        ```
        >>> spark.createDataFrame([[1]], "a int").show()
        ...
        ValueError: Could not parse datatype: a int
        ```
      
        **After**
        ```
        >>> spark.createDataFrame([[[1]]], "struct<a: struct<b: int>>").show()
        +---+
        |  a|
        +---+
        |[1]|
        +---+
        ```
      
        ```
        >>> spark.createDataFrame([[[1]]], "a: struct<b: int>").show()
        +---+
        |  a|
        +---+
        |[1]|
        +---+
        ```
      
        ```
        >>> spark.createDataFrame([[1]], "a int").show()
        +---+
        |  a|
        +---+
        |  1|
        +---+
        ```
      
      ## How was this patch tested?
      
      Author: hyukjinkwon <gurwls223@gmail.com>
      
      Closes #18590 from HyukjinKwon/deduplicate-python-ddl.
      ebc124d4
    • Xingbo Jiang's avatar
      [SPARK-21366][SQL][TEST] Add sql test for window functions · 66d21686
      Xingbo Jiang authored
      ## What changes were proposed in this pull request?
      
      Add sql test for window functions, also remove uncecessary test cases in `WindowQuerySuite`.
      
      ## How was this patch tested?
      
      Added `window.sql` and the corresponding output file.
      
      Author: Xingbo Jiang <xingbo.jiang@databricks.com>
      
      Closes #18591 from jiangxb1987/window.
      66d21686
    • hyukjinkwon's avatar
      [SPARK-21263][SQL] Do not allow partially parsing double and floats via NumberFormat in CSV · 7514db1d
      hyukjinkwon authored
      ## What changes were proposed in this pull request?
      
      This PR proposes to remove `NumberFormat.parse` use to disallow a case of partially parsed data. For example,
      
      ```
      scala> spark.read.schema("a DOUBLE").option("mode", "FAILFAST").csv(Seq("10u12").toDS).show()
      +----+
      |   a|
      +----+
      |10.0|
      +----+
      ```
      
      ## How was this patch tested?
      
      Unit tests added in `UnivocityParserSuite` and `CSVSuite`.
      
      Author: hyukjinkwon <gurwls223@gmail.com>
      
      Closes #18532 from HyukjinKwon/SPARK-21263.
      7514db1d
  9. Jul 10, 2017
    • jinxing's avatar
      [SPARK-21315][SQL] Skip some spill files when generateIterator(startIndex) in... · 97a1aa2c
      jinxing authored
      [SPARK-21315][SQL] Skip some spill files when generateIterator(startIndex) in ExternalAppendOnlyUnsafeRowArray.
      
      ## What changes were proposed in this pull request?
      
      In current code, it is expensive to use `UnboundedFollowingWindowFunctionFrame`, because it is iterating from the start to lower bound every time calling `write` method. When traverse the iterator, it's possible to skip some spilled files thus to save some time.
      
      ## How was this patch tested?
      
      Added unit test
      
      Did a small test for benchmark:
      
      Put 2000200 rows into `UnsafeExternalSorter`-- 2 spill files(each contains 1000000 rows) and inMemSorter contains 200 rows.
      Move the iterator forward to index=2000001.
      
      *With this change*:
      `getIterator(2000001)`, it will cost almost 0ms~1ms;
      *Without this change*:
      `for(int i=0; i<2000001; i++)geIterator().loadNext()`, it will cost 300ms.
      
      Author: jinxing <jinxing6042@126.com>
      
      Closes #18541 from jinxing64/SPARK-21315.
      97a1aa2c
    • gatorsmile's avatar
      [SPARK-21350][SQL] Fix the error message when the number of arguments is wrong when invoking a UDF · 1471ee7a
      gatorsmile authored
      ### What changes were proposed in this pull request?
      Users get a very confusing error when users specify a wrong number of parameters.
      ```Scala
          val df = spark.emptyDataFrame
          spark.udf.register("foo", (_: String).length)
          df.selectExpr("foo(2, 3, 4)")
      ```
      ```
      org.apache.spark.sql.UDFSuite$$anonfun$9$$anonfun$apply$mcV$sp$12 cannot be cast to scala.Function3
      java.lang.ClassCastException: org.apache.spark.sql.UDFSuite$$anonfun$9$$anonfun$apply$mcV$sp$12 cannot be cast to scala.Function3
      	at org.apache.spark.sql.catalyst.expressions.ScalaUDF.<init>(ScalaUDF.scala:109)
      ```
      
      This PR is to capture the exception and issue an error message that is consistent with what we did for built-in functions. After the fix, the error message is improved to
      ```
      Invalid number of arguments for function foo; line 1 pos 0
      org.apache.spark.sql.AnalysisException: Invalid number of arguments for function foo; line 1 pos 0
      	at org.apache.spark.sql.catalyst.analysis.SimpleFunctionRegistry.lookupFunction(FunctionRegistry.scala:119)
      ```
      
      ### How was this patch tested?
      Added a test case
      
      Author: gatorsmile <gatorsmile@gmail.com>
      
      Closes #18574 from gatorsmile/statsCheck.
      1471ee7a
    • Takeshi Yamamuro's avatar
      [SPARK-21043][SQL] Add unionByName in Dataset · a2bec6c9
      Takeshi Yamamuro authored
      ## What changes were proposed in this pull request?
      This pr added `unionByName` in `DataSet`.
      Here is how to use:
      ```
      val df1 = Seq((1, 2, 3)).toDF("col0", "col1", "col2")
      val df2 = Seq((4, 5, 6)).toDF("col1", "col2", "col0")
      df1.unionByName(df2).show
      
      // output:
      // +----+----+----+
      // |col0|col1|col2|
      // +----+----+----+
      // |   1|   2|   3|
      // |   6|   4|   5|
      // +----+----+----+
      ```
      
      ## How was this patch tested?
      Added tests in `DataFrameSuite`.
      
      Author: Takeshi Yamamuro <yamamuro@apache.org>
      
      Closes #18300 from maropu/SPARK-21043-2.
      a2bec6c9
    • Bryan Cutler's avatar
      [SPARK-13534][PYSPARK] Using Apache Arrow to increase performance of DataFrame.toPandas · d03aebbe
      Bryan Cutler authored
      ## What changes were proposed in this pull request?
      Integrate Apache Arrow with Spark to increase performance of `DataFrame.toPandas`.  This has been done by using Arrow to convert data partitions on the executor JVM to Arrow payload byte arrays where they are then served to the Python process.  The Python DataFrame can then collect the Arrow payloads where they are combined and converted to a Pandas DataFrame.  Data types except complex, date, timestamp, and decimal  are currently supported, otherwise an `UnsupportedOperation` exception is thrown.
      
      Additions to Spark include a Scala package private method `Dataset.toArrowPayload` that will convert data partitions in the executor JVM to `ArrowPayload`s as byte arrays so they can be easily served.  A package private class/object `ArrowConverters` that provide data type mappings and conversion routines.  In Python, a private method `DataFrame._collectAsArrow` is added to collect Arrow payloads and a SQLConf "spark.sql.execution.arrow.enable" can be used in `toPandas()` to enable using Arrow (uses the old conversion by default).
      
      ## How was this patch tested?
      Added a new test suite `ArrowConvertersSuite` that will run tests on conversion of Datasets to Arrow payloads for supported types.  The suite will generate a Dataset and matching Arrow JSON data, then the dataset is converted to an Arrow payload and finally validated against the JSON data.  This will ensure that the schema and data has been converted correctly.
      
      Added PySpark tests to verify the `toPandas` method is producing equal DataFrames with and without pyarrow.  A roundtrip test to ensure the pandas DataFrame produced by pyspark is equal to a one made directly with pandas.
      
      Author: Bryan Cutler <cutlerb@gmail.com>
      Author: Li Jin <ice.xelloss@gmail.com>
      Author: Li Jin <li.jin@twosigma.com>
      Author: Wes McKinney <wes.mckinney@twosigma.com>
      
      Closes #18459 from BryanCutler/toPandas_with_arrow-SPARK-13534.
      d03aebbe
    • hyukjinkwon's avatar
      [SPARK-21266][R][PYTHON] Support schema a DDL-formatted string in dapply/gapply/from_json · 2bfd5acc
      hyukjinkwon authored
      ## What changes were proposed in this pull request?
      
      This PR supports schema in a DDL formatted string for `from_json` in R/Python and `dapply` and `gapply` in R, which are commonly used and/or consistent with Scala APIs.
      
      Additionally, this PR exposes `structType` in R to allow working around in other possible corner cases.
      
      **Python**
      
      `from_json`
      
      ```python
      from pyspark.sql.functions import from_json
      
      data = [(1, '''{"a": 1}''')]
      df = spark.createDataFrame(data, ("key", "value"))
      df.select(from_json(df.value, "a INT").alias("json")).show()
      ```
      
      **R**
      
      `from_json`
      
      ```R
      df <- sql("SELECT named_struct('name', 'Bob') as people")
      df <- mutate(df, people_json = to_json(df$people))
      head(select(df, from_json(df$people_json, "name STRING")))
      ```
      
      `structType.character`
      
      ```R
      structType("a STRING, b INT")
      ```
      
      `dapply`
      
      ```R
      dapply(createDataFrame(list(list(1.0)), "a"), function(x) {x}, "a DOUBLE")
      ```
      
      `gapply`
      
      ```R
      gapply(createDataFrame(list(list(1.0)), "a"), "a", function(key, x) { x }, "a DOUBLE")
      ```
      
      ## How was this patch tested?
      
      Doc tests for `from_json` in Python and unit tests `test_sparkSQL.R` in R.
      
      Author: hyukjinkwon <gurwls223@gmail.com>
      
      Closes #18498 from HyukjinKwon/SPARK-21266.
      2bfd5acc
    • Juliusz Sompolski's avatar
      [SPARK-21272] SortMergeJoin LeftAnti does not update numOutputRows · 18b3b00e
      Juliusz Sompolski authored
      ## What changes were proposed in this pull request?
      
      Updating numOutputRows metric was missing from one return path of LeftAnti SortMergeJoin.
      
      ## How was this patch tested?
      
      Non-zero output rows manually seen in metrics.
      
      Author: Juliusz Sompolski <julek@databricks.com>
      
      Closes #18494 from juliuszsompolski/SPARK-21272.
      18b3b00e
    • Takeshi Yamamuro's avatar
      [SPARK-20460][SQL] Make it more consistent to handle column name duplication · 647963a2
      Takeshi Yamamuro authored
      ## What changes were proposed in this pull request?
      This pr made it more consistent to handle column name duplication. In the current master, error handling is different when hitting column name duplication:
      ```
      // json
      scala> val schema = StructType(StructField("a", IntegerType) :: StructField("a", IntegerType) :: Nil)
      scala> Seq("""{"a":1, "a":1}"""""").toDF().coalesce(1).write.mode("overwrite").text("/tmp/data")
      scala> spark.read.format("json").schema(schema).load("/tmp/data").show
      org.apache.spark.sql.AnalysisException: Reference 'a' is ambiguous, could be: a#12, a#13.;
        at org.apache.spark.sql.catalyst.plans.logical.LogicalPlan.resolve(LogicalPlan.scala:287)
        at org.apache.spark.sql.catalyst.plans.logical.LogicalPlan.resolve(LogicalPlan.scala:181)
        at org.apache.spark.sql.catalyst.plans.logical.LogicalPlan$$anonfun$resolve$1.apply(LogicalPlan.scala:153)
      
      scala> spark.read.format("json").load("/tmp/data").show
      org.apache.spark.sql.AnalysisException: Duplicate column(s) : "a" found, cannot save to JSON format;
        at org.apache.spark.sql.execution.datasources.json.JsonDataSource.checkConstraints(JsonDataSource.scala:81)
        at org.apache.spark.sql.execution.datasources.json.JsonDataSource.inferSchema(JsonDataSource.scala:63)
        at org.apache.spark.sql.execution.datasources.json.JsonFileFormat.inferSchema(JsonFileFormat.scala:57)
        at org.apache.spark.sql.execution.datasources.DataSource$$anonfun$7.apply(DataSource.scala:176)
        at org.apache.spark.sql.execution.datasources.DataSource$$anonfun$7.apply(DataSource.scala:176)
      
      // csv
      scala> val schema = StructType(StructField("a", IntegerType) :: StructField("a", IntegerType) :: Nil)
      scala> Seq("a,a", "1,1").toDF().coalesce(1).write.mode("overwrite").text("/tmp/data")
      scala> spark.read.format("csv").schema(schema).option("header", false).load("/tmp/data").show
      org.apache.spark.sql.AnalysisException: Reference 'a' is ambiguous, could be: a#41, a#42.;
        at org.apache.spark.sql.catalyst.plans.logical.LogicalPlan.resolve(LogicalPlan.scala:287)
        at org.apache.spark.sql.catalyst.plans.logical.LogicalPlan.resolve(LogicalPlan.scala:181)
        at org.apache.spark.sql.catalyst.plans.logical.LogicalPlan$$anonfun$resolve$1.apply(LogicalPlan.scala:153)
        at org.apache.spark.sql.catalyst.plans.logical.LogicalPlan$$anonfun$resolve$1.apply(LogicalPlan.scala:152)
      
      // If `inferSchema` is true, a CSV format is duplicate-safe (See SPARK-16896)
      scala> spark.read.format("csv").option("header", true).load("/tmp/data").show
      +---+---+
      | a0| a1|
      +---+---+
      |  1|  1|
      +---+---+
      
      // parquet
      scala> val schema = StructType(StructField("a", IntegerType) :: StructField("a", IntegerType) :: Nil)
      scala> Seq((1, 1)).toDF("a", "b").coalesce(1).write.mode("overwrite").parquet("/tmp/data")
      scala> spark.read.format("parquet").schema(schema).option("header", false).load("/tmp/data").show
      org.apache.spark.sql.AnalysisException: Reference 'a' is ambiguous, could be: a#110, a#111.;
        at org.apache.spark.sql.catalyst.plans.logical.LogicalPlan.resolve(LogicalPlan.scala:287)
        at org.apache.spark.sql.catalyst.plans.logical.LogicalPlan.resolve(LogicalPlan.scala:181)
        at org.apache.spark.sql.catalyst.plans.logical.LogicalPlan$$anonfun$resolve$1.apply(LogicalPlan.scala:153)
        at org.apache.spark.sql.catalyst.plans.logical.LogicalPlan$$anonfun$resolve$1.apply(LogicalPlan.scala:152)
        at scala.collection.TraversableLike$$anonfun$map$1.apply(TraversableLike.scala:234)
        at scala.collection.TraversableLike$$anonfun$map$1.apply(TraversableLike.scala:234)
      ```
      When this patch applied, the results change to;
      ```
      
      // json
      scala> val schema = StructType(StructField("a", IntegerType) :: StructField("a", IntegerType) :: Nil)
      scala> Seq("""{"a":1, "a":1}"""""").toDF().coalesce(1).write.mode("overwrite").text("/tmp/data")
      scala> spark.read.format("json").schema(schema).load("/tmp/data").show
      org.apache.spark.sql.AnalysisException: Found duplicate column(s) in datasource: "a";
        at org.apache.spark.sql.util.SchemaUtils$.checkColumnNameDuplication(SchemaUtil.scala:47)
        at org.apache.spark.sql.util.SchemaUtils$.checkSchemaColumnNameDuplication(SchemaUtil.scala:33)
        at org.apache.spark.sql.execution.datasources.DataSource.getOrInferFileFormatSchema(DataSource.scala:186)
        at org.apache.spark.sql.execution.datasources.DataSource.resolveRelation(DataSource.scala:368)
      
      scala> spark.read.format("json").load("/tmp/data").show
      org.apache.spark.sql.AnalysisException: Found duplicate column(s) in datasource: "a";
        at org.apache.spark.sql.util.SchemaUtils$.checkColumnNameDuplication(SchemaUtil.scala:47)
        at org.apache.spark.sql.util.SchemaUtils$.checkSchemaColumnNameDuplication(SchemaUtil.scala:33)
        at org.apache.spark.sql.execution.datasources.DataSource.getOrInferFileFormatSchema(DataSource.scala:186)
        at org.apache.spark.sql.execution.datasources.DataSource.resolveRelation(DataSource.scala:368)
        at org.apache.spark.sql.DataFrameReader.load(DataFrameReader.scala:178)
        at org.apache.spark.sql.DataFrameReader.load(DataFrameReader.scala:156)
      
      // csv
      scala> val schema = StructType(StructField("a", IntegerType) :: StructField("a", IntegerType) :: Nil)
      scala> Seq("a,a", "1,1").toDF().coalesce(1).write.mode("overwrite").text("/tmp/data")
      scala> spark.read.format("csv").schema(schema).option("header", false).load("/tmp/data").show
      org.apache.spark.sql.AnalysisException: Found duplicate column(s) in datasource: "a";
        at org.apache.spark.sql.util.SchemaUtils$.checkColumnNameDuplication(SchemaUtil.scala:47)
        at org.apache.spark.sql.util.SchemaUtils$.checkSchemaColumnNameDuplication(SchemaUtil.scala:33)
        at org.apache.spark.sql.execution.datasources.DataSource.getOrInferFileFormatSchema(DataSource.scala:186)
        at org.apache.spark.sql.execution.datasources.DataSource.resolveRelation(DataSource.scala:368)
        at org.apache.spark.sql.DataFrameReader.load(DataFrameReader.scala:178)
      
      scala> spark.read.format("csv").option("header", true).load("/tmp/data").show
      +---+---+
      | a0| a1|
      +---+---+
      |  1|  1|
      +---+---+
      
      // parquet
      scala> val schema = StructType(StructField("a", IntegerType) :: StructField("a", IntegerType) :: Nil)
      scala> Seq((1, 1)).toDF("a", "b").coalesce(1).write.mode("overwrite").parquet("/tmp/data")
      scala> spark.read.format("parquet").schema(schema).option("header", false).load("/tmp/data").show
      org.apache.spark.sql.AnalysisException: Found duplicate column(s) in datasource: "a";
        at org.apache.spark.sql.util.SchemaUtils$.checkColumnNameDuplication(SchemaUtil.scala:47)
        at org.apache.spark.sql.util.SchemaUtils$.checkSchemaColumnNameDuplication(SchemaUtil.scala:33)
        at org.apache.spark.sql.execution.datasources.DataSource.getOrInferFileFormatSchema(DataSource.scala:186)
        at org.apache.spark.sql.execution.datasources.DataSource.resolveRelation(DataSource.scala:368)
      ```
      
      ## How was this patch tested?
      Added tests in `DataFrameReaderWriterSuite` and `SQLQueryTestSuite`.
      
      Author: Takeshi Yamamuro <yamamuro@apache.org>
      
      Closes #17758 from maropu/SPARK-20460.
      647963a2
    • Wenchen Fan's avatar
      [SPARK-21100][SQL][FOLLOWUP] cleanup code and add more comments for Dataset.summary · 0e80ecae
      Wenchen Fan authored
      ## What changes were proposed in this pull request?
      
      Some code cleanup and adding comments to make the code more readable. Changed the way to generate result rows, to be more clear.
      
      ## How was this patch tested?
      
      existing tests
      
      Author: Wenchen Fan <wenchen@databricks.com>
      
      Closes #18570 from cloud-fan/summary.
      0e80ecae
  10. Jul 09, 2017
    • Wenchen Fan's avatar
      [SPARK-18016][SQL][FOLLOWUP] merge declareAddedFunctions, initNestedClasses... · 680b33f1
      Wenchen Fan authored
      [SPARK-18016][SQL][FOLLOWUP] merge declareAddedFunctions, initNestedClasses and declareNestedClasses
      
      ## What changes were proposed in this pull request?
      
      These 3 methods have to be used together, so it makes more sense to merge them into one method and then the caller side only need to call one method.
      
      ## How was this patch tested?
      
      existing tests.
      
      Author: Wenchen Fan <wenchen@databricks.com>
      
      Closes #18579 from cloud-fan/minor.
      680b33f1
  11. Jul 08, 2017
    • Xiao Li's avatar
      [SPARK-21307][REVERT][SQL] Remove SQLConf parameters from the parser-related classes · c3712b77
      Xiao Li authored
      ## What changes were proposed in this pull request?
      Since we do not set active sessions when parsing the plan, we are unable to correctly use SQLConf.get to find the correct active session. Since https://github.com/apache/spark/pull/18531 breaks the build, I plan to revert it at first.
      
      ## How was this patch tested?
      The existing test cases
      
      Author: Xiao Li <gatorsmile@gmail.com>
      
      Closes #18568 from gatorsmile/revert18531.
      c3712b77
    • Zhenhua Wang's avatar
      [SPARK-21083][SQL] Store zero size and row count when analyzing empty table · 9fccc362
      Zhenhua Wang authored
      ## What changes were proposed in this pull request?
      
      We should be able to store zero size and row count after analyzing empty table.
      
      This pr also enhances the test cases for re-analyzing tables.
      
      ## How was this patch tested?
      
      Added a new test case and enhanced some test cases.
      
      Author: Zhenhua Wang <wangzhenhua@huawei.com>
      
      Closes #18292 from wzhfy/analyzeNewColumn.
      9fccc362
    • Dongjoon Hyun's avatar
      [SPARK-21345][SQL][TEST][TEST-MAVEN] SparkSessionBuilderSuite should clean up stopped sessions. · 0b8dd2d0
      Dongjoon Hyun authored
      ## What changes were proposed in this pull request?
      
      `SparkSessionBuilderSuite` should clean up stopped sessions. Otherwise, it leaves behind some stopped `SparkContext`s interfereing with other test suites using `ShardSQLContext`.
      
      Recently, master branch fails consequtively.
      - https://amplab.cs.berkeley.edu/jenkins/view/Spark%20QA%20Test%20(Dashboard)/
      
      ## How was this patch tested?
      
      Pass the Jenkins with a updated suite.
      
      Author: Dongjoon Hyun <dongjoon@apache.org>
      
      Closes #18567 from dongjoon-hyun/SPARK-SESSION.
      0b8dd2d0
    • Michael Patterson's avatar
      [SPARK-20456][DOCS] Add examples for functions collection for pyspark · f5f02d21
      Michael Patterson authored
      ## What changes were proposed in this pull request?
      
      This adds documentation to many functions in pyspark.sql.functions.py:
      `upper`, `lower`, `reverse`, `unix_timestamp`, `from_unixtime`, `rand`, `randn`, `collect_list`, `collect_set`, `lit`
      Add units to the trigonometry functions.
      Renames columns in datetime examples to be more informative.
      Adds links between some functions.
      
      ## How was this patch tested?
      
      `./dev/lint-python`
      `python python/pyspark/sql/functions.py`
      `./python/run-tests.py --module pyspark-sql`
      
      Author: Michael Patterson <map222@gmail.com>
      
      Closes #17865 from map222/spark-20456.
      f5f02d21
    • Takeshi Yamamuro's avatar
      [SPARK-21281][SQL] Use string types by default if array and map have no argument · 7896e7b9
      Takeshi Yamamuro authored
      ## What changes were proposed in this pull request?
      This pr modified code to use string types by default if `array` and `map` in functions have no argument. This behaviour is the same with Hive one;
      ```
      hive> CREATE TEMPORARY TABLE t1 AS SELECT map();
      hive> DESCRIBE t1;
      _c0   map<string,string>
      
      hive> CREATE TEMPORARY TABLE t2 AS SELECT array();
      hive> DESCRIBE t2;
      _c0   array<string>
      ```
      
      ## How was this patch tested?
      Added tests in `DataFrameFunctionsSuite`.
      
      Author: Takeshi Yamamuro <yamamuro@apache.org>
      
      Closes #18516 from maropu/SPARK-21281.
      7896e7b9
    • Andrew Ray's avatar
      [SPARK-21100][SQL] Add summary method as alternative to describe that gives... · e1a172c2
      Andrew Ray authored
      [SPARK-21100][SQL] Add summary method as alternative to describe that gives quartiles similar to Pandas
      
      ## What changes were proposed in this pull request?
      
      Adds method `summary`  that allows user to specify which statistics and percentiles to calculate. By default it include the existing statistics from `describe` and quartiles (25th, 50th, and 75th percentiles) similar to Pandas. Also changes the implementation of `describe` to delegate to `summary`.
      
      ## How was this patch tested?
      
      additional unit test
      
      Author: Andrew Ray <ray.andrew@gmail.com>
      
      Closes #18307 from aray/SPARK-21100.
      e1a172c2
  12. Jul 07, 2017
    • Wang Gengliang's avatar
      [SPARK-21336] Revise rand comparison in BatchEvalPythonExecSuite · a0fe32a2
      Wang Gengliang authored
      ## What changes were proposed in this pull request?
      
      Revise rand comparison in BatchEvalPythonExecSuite
      
      In BatchEvalPythonExecSuite, there are two cases using the case "rand() > 3"
      Rand() generates a random value in [0, 1), it is wired to be compared with 3, use 0.3 instead
      
      ## How was this patch tested?
      
      unit test
      
      Please review http://spark.apache.org/contributing.html before opening a pull request.
      
      Author: Wang Gengliang <ltnwgl@gmail.com>
      
      Closes #18560 from gengliangwang/revise_BatchEvalPythonExecSuite.
      a0fe32a2
    • Wenchen Fan's avatar
      [SPARK-21335][SQL] support un-aliased subquery · fef08130
      Wenchen Fan authored
      ## What changes were proposed in this pull request?
      
      un-aliased subquery is supported by Spark SQL for a long time. Its semantic was not well defined and had confusing behaviors, and it's not a standard SQL syntax, so we disallowed it in https://issues.apache.org/jira/browse/SPARK-20690 .
      
      However, this is a breaking change, and we do have existing queries using un-aliased subquery. We should add the support back and fix its semantic.
      
      This PR fixes the un-aliased subquery by assigning a default alias name.
      
      After this PR, there is no syntax change from branch 2.2 to master, but we invalid a weird use case:
      `SELECT v.i from (SELECT i FROM v)`. Now this query will throw analysis exception because users should not be able to use the qualifier inside a subquery.
      
      ## How was this patch tested?
      
      new regression test
      
      Author: Wenchen Fan <wenchen@databricks.com>
      
      Closes #18559 from cloud-fan/sub-query.
      fef08130
    • Jacek Laskowski's avatar
      [SPARK-21313][SS] ConsoleSink's string representation · 7fcbb9b5
      Jacek Laskowski authored
      ## What changes were proposed in this pull request?
      
      Add `toString` with options for `ConsoleSink` so it shows nicely in query progress.
      
      **BEFORE**
      
      ```
        "sink" : {
          "description" : "org.apache.spark.sql.execution.streaming.ConsoleSink4b340441"
        }
      ```
      
      **AFTER**
      
      ```
        "sink" : {
          "description" : "ConsoleSink[numRows=10, truncate=false]"
        }
      ```
      
      /cc zsxwing tdas
      
      ## How was this patch tested?
      
      Local build
      
      Author: Jacek Laskowski <jacek@japila.pl>
      
      Closes #18539 from jaceklaskowski/SPARK-21313-ConsoleSink-toString.
      7fcbb9b5
    • Liang-Chi Hsieh's avatar
      [SPARK-20703][SQL][FOLLOW-UP] Associate metrics with data writes onto DataFrameWriter operations · 5df99bd3
      Liang-Chi Hsieh authored
      ## What changes were proposed in this pull request?
      
      Remove time metrics since it seems no way to measure it in non per-row tracking.
      
      ## How was this patch tested?
      
      Existing tests.
      
      Please review http://spark.apache.org/contributing.html before opening a pull request.
      
      Author: Liang-Chi Hsieh <viirya@gmail.com>
      
      Closes #18558 from viirya/SPARK-20703-followup.
      5df99bd3
    • Kazuaki Ishizaki's avatar
      [SPARK-21217][SQL] Support ColumnVector.Array.to<type>Array() · c09b31eb
      Kazuaki Ishizaki authored
      ## What changes were proposed in this pull request?
      
      This PR implements bulk-copy for `ColumnVector.Array.to<type>Array()` methods (e.g. `toIntArray()`) in `ColumnVector.Array` by using `System.arrayCopy()` or `Platform.copyMemory()`.
      
      Before this PR, when one of these method is called, the generic method in `ArrayData` is called. It is not fast since element-wise copy is performed.
      
      This PR can improve performance of a benchmark program by 1.9x and 3.2x.
      
      Without this PR
      ```
      OpenJDK 64-Bit Server VM 1.8.0_131-8u131-b11-0ubuntu1.16.04.2-b11 on Linux 4.4.0-66-generic
      Intel(R) Xeon(R) CPU E5-2667 v3  3.20GHz
      
      Int Array                                Best/Avg Time(ms)    Rate(M/s)   Per Row(ns)
      ------------------------------------------------------------------------------------------------
      ON_HEAP                                        586 /  628         14.3          69.9
      OFF_HEAP                                       893 /  902          9.4         106.5
      ```
      
      With this PR
      ```
      OpenJDK 64-Bit Server VM 1.8.0_131-8u131-b11-0ubuntu1.16.04.2-b11 on Linux 4.4.0-66-generic
      Intel(R) Xeon(R) CPU E5-2667 v3  3.20GHz
      
      Int Array                                Best/Avg Time(ms)    Rate(M/s)   Per Row(ns)
      ------------------------------------------------------------------------------------------------
      ON_HEAP                                        306 /  331         27.4          36.4
      OFF_HEAP                                       282 /  287         29.8          33.6
      ```
      
      Source program
      ```
          (MemoryMode.ON_HEAP :: MemoryMode.OFF_HEAP :: Nil).foreach { memMode => {
            val len = 8 * 1024 * 1024
            val column = ColumnVector.allocate(len * 2, new ArrayType(IntegerType, false), memMode)
      
            val data = column.arrayData
            var i = 0
            while (i < len) {
              data.putInt(i, i)
              i += 1
            }
            column.putArray(0, 0, len)
      
            val benchmark = new Benchmark("Int Array", len, minNumIters = 20)
            benchmark.addCase(s"$memMode") { iter =>
              var i = 0
              while (i < 50) {
                column.getArray(0).toIntArray
                i += 1
              }
            }
            benchmark.run
          }}
      ```
      
      ## How was this patch tested?
      
      Added test suite
      
      Author: Kazuaki Ishizaki <ishizaki@jp.ibm.com>
      
      Closes #18425 from kiszk/SPARK-21217.
      c09b31eb
  13. Jul 06, 2017
Loading