Skip to content
Snippets Groups Projects
  1. Nov 02, 2016
    • Steve Loughran's avatar
      [SPARK-17058][BUILD] Add maven snapshots-and-staging profile to build/test... · 1eef8e5c
      Steve Loughran authored
      [SPARK-17058][BUILD] Add maven snapshots-and-staging profile to build/test against staging artifacts
      
      ## What changes were proposed in this pull request?
      
      Adds a `snapshots-and-staging profile` so that  RCs of projects like Hadoop and HBase can be used in developer-only build and test runs. There's a comment above the profile telling people not to use this in production.
      
      There's no attempt to do the same for SBT, as Ivy is different.
      ## How was this patch tested?
      
      Tested by building against the Hadoop 2.7.3 RC 1 JARs
      
      without the profile (and without any local copy of the 2.7.3 artifacts), the build failed
      
      ```
      mvn install -DskipTests -Pyarn,hadoop-2.7,hive -Dhadoop.version=2.7.3
      
      ...
      
      [INFO] ------------------------------------------------------------------------
      [INFO] Building Spark Project Launcher 2.1.0-SNAPSHOT
      [INFO] ------------------------------------------------------------------------
      Downloading: https://repo1.maven.org/maven2/org/apache/hadoop/hadoop-client/2.7.3/hadoop-client-2.7.3.pom
      [WARNING] The POM for org.apache.hadoop:hadoop-client:jar:2.7.3 is missing, no dependency information available
      Downloading: https://repo1.maven.org/maven2/org/apache/hadoop/hadoop-client/2.7.3/hadoop-client-2.7.3.jar
      
      
      [INFO] ------------------------------------------------------------------------
      [INFO] Reactor Summary:
      [INFO]
      [INFO] Spark Project Parent POM ........................... SUCCESS [  4.482 s]
      [INFO] Spark Project Tags ................................. SUCCESS [ 17.402 s]
      [INFO] Spark Project Sketch ............................... SUCCESS [ 11.252 s]
      [INFO] Spark Project Networking ........................... SUCCESS [ 13.458 s]
      [INFO] Spark Project Shuffle Streaming Service ............ SUCCESS [  9.043 s]
      [INFO] Spark Project Unsafe ............................... SUCCESS [ 16.027 s]
      [INFO] Spark Project Launcher ............................. FAILURE [  1.653 s]
      [INFO] Spark Project Core ................................. SKIPPED
      ...
      ```
      
      With the profile, the build completed
      
      ```
      mvn install -DskipTests -Pyarn,hadoop-2.7,hive,snapshots-and-staging -Dhadoop.version=2.7.3
      ```
      
      Author: Steve Loughran <stevel@apache.org>
      
      Closes #14646 from steveloughran/stevel/SPARK-17058-support-asf-snapshots.
      
      (cherry picked from commit 37d95227)
      Signed-off-by: default avatarReynold Xin <rxin@databricks.com>
      1eef8e5c
    • Jeff Zhang's avatar
      [SPARK-18160][CORE][YARN] spark.files & spark.jars should not be passed to driver in yarn mode · bd3ea659
      Jeff Zhang authored
      
      ## What changes were proposed in this pull request?
      
      spark.files is still passed to driver in yarn mode, so SparkContext will still handle it which cause the error in the jira desc.
      
      ## How was this patch tested?
      
      Tested manually in a 5 node cluster. As this issue only happens in multiple node cluster, so I didn't write test for it.
      
      Author: Jeff Zhang <zjffdu@apache.org>
      
      Closes #15669 from zjffdu/SPARK-18160.
      
      (cherry picked from commit 3c24299b)
      Signed-off-by: default avatarMarcelo Vanzin <vanzin@cloudera.com>
      bd3ea659
    • Xiangrui Meng's avatar
      [SPARK-14393][SQL] values generated by non-deterministic functions shouldn't... · 0093257e
      Xiangrui Meng authored
      [SPARK-14393][SQL] values generated by non-deterministic functions shouldn't change after coalesce or union
      
      ## What changes were proposed in this pull request?
      
      When a user appended a column using a "nondeterministic" function to a DataFrame, e.g., `rand`, `randn`, and `monotonically_increasing_id`, the expected semantic is the following:
      - The value in each row should remain unchanged, as if we materialize the column immediately, regardless of later DataFrame operations.
      
      However, since we use `TaskContext.getPartitionId` to get the partition index from the current thread, the values from nondeterministic columns might change if we call `union` or `coalesce` after. `TaskContext.getPartitionId` returns the partition index of the current Spark task, which might not be the corresponding partition index of the DataFrame where we defined the column.
      
      See the unit tests below or JIRA for examples.
      
      This PR uses the partition index from `RDD.mapPartitionWithIndex` instead of `TaskContext` and fixes the partition initialization logic in whole-stage codegen, normal codegen, and codegen fallback. `initializeStatesForPartition(partitionIndex: Int)` was added to `Projection`, `Nondeterministic`, and `Predicate` (codegen) and initialized right after object creation in `mapPartitionWithIndex`. `newPredicate` now returns a `Predicate` instance rather than a function for proper initialization.
      ## How was this patch tested?
      
      Unit tests. (Actually I'm not very confident that this PR fixed all issues without introducing new ones ...)
      
      cc: rxin davies
      
      Author: Xiangrui Meng <meng@databricks.com>
      
      Closes #15567 from mengxr/SPARK-14393.
      
      (cherry picked from commit 02f20310)
      Signed-off-by: default avatarReynold Xin <rxin@databricks.com>
      0093257e
    • buzhihuojie's avatar
      [SPARK-17895] Improve doc for rangeBetween and rowsBetween · a885d5bb
      buzhihuojie authored
      ## What changes were proposed in this pull request?
      
      Copied description for row and range based frame boundary from https://github.com/apache/spark/blob/master/sql/core/src/main/scala/org/apache/spark/sql/execution/window/WindowExec.scala#L56
      
      Added examples to show different behavior of rangeBetween and rowsBetween when involving duplicate values.
      
      Please review https://cwiki.apache.org/confluence/display/SPARK/Contributing+to+Spark
      
       before opening a pull request.
      
      Author: buzhihuojie <ren.weiluo@gmail.com>
      
      Closes #15727 from david-weiluo-ren/improveDocForRangeAndRowsBetween.
      
      (cherry picked from commit 742e0fea)
      Signed-off-by: default avatarReynold Xin <rxin@databricks.com>
      a885d5bb
    • Takeshi YAMAMURO's avatar
      [SPARK-17683][SQL] Support ArrayType in Literal.apply · 9be06912
      Takeshi YAMAMURO authored
      
      ## What changes were proposed in this pull request?
      
      This pr is to add pattern-matching entries for array data in `Literal.apply`.
      ## How was this patch tested?
      
      Added tests in `LiteralExpressionSuite`.
      
      Author: Takeshi YAMAMURO <linguin.m.s@gmail.com>
      
      Closes #15257 from maropu/SPARK-17683.
      
      (cherry picked from commit 4af0ce2d)
      Signed-off-by: default avatarReynold Xin <rxin@databricks.com>
      9be06912
    • eyal farago's avatar
      [SPARK-16839][SQL] Simplify Struct creation code path · 41491e54
      eyal farago authored
      
      ## What changes were proposed in this pull request?
      
      Simplify struct creation, especially the aspect of `CleanupAliases` which missed some aliases when handling trees created by `CreateStruct`.
      
      This PR includes:
      
      1. A failing test (create struct with nested aliases, some of the aliases survive `CleanupAliases`).
      2. A fix that transforms `CreateStruct` into a `CreateNamedStruct` constructor, effectively eliminating `CreateStruct` from all expression trees.
      3. A `NamePlaceHolder` used by `CreateStruct` when column names cannot be extracted from unresolved `NamedExpression`.
      4. A new Analyzer rule that resolves `NamePlaceHolder` into a string literal once the `NamedExpression` is resolved.
      5. `CleanupAliases` code was simplified as it no longer has to deal with `CreateStruct`'s top level columns.
      
      ## How was this patch tested?
      Running all tests-suits in package org.apache.spark.sql, especially including the analysis suite, making sure added test initially fails, after applying suggested fix rerun the entire analysis package successfully.
      
      Modified few tests that expected `CreateStruct` which is now transformed into `CreateNamedStruct`.
      
      Author: eyal farago <eyal farago>
      Author: Herman van Hovell <hvanhovell@databricks.com>
      Author: eyal farago <eyal.farago@gmail.com>
      Author: Eyal Farago <eyal.farago@actimize.com>
      Author: Hyukjin Kwon <gurwls223@gmail.com>
      Author: eyalfa <eyal.farago@gmail.com>
      
      Closes #15718 from hvanhovell/SPARK-16839-2.
      
      (cherry picked from commit f151bd1a)
      Signed-off-by: default avatarHerman van Hovell <hvanhovell@databricks.com>
      41491e54
    • Sean Owen's avatar
      [SPARK-18076][CORE][SQL] Fix default Locale used in DateFormat, NumberFormat to Locale.US · 176afa5e
      Sean Owen authored
      
      ## What changes were proposed in this pull request?
      
      Fix `Locale.US` for all usages of `DateFormat`, `NumberFormat`
      ## How was this patch tested?
      
      Existing tests.
      
      Author: Sean Owen <sowen@cloudera.com>
      
      Closes #15610 from srowen/SPARK-18076.
      
      (cherry picked from commit 9c8deef6)
      Signed-off-by: default avatarSean Owen <sowen@cloudera.com>
      Unverified
      176afa5e
    • Liwei Lin's avatar
      [SPARK-18198][DOC][STREAMING] Highlight code snippets · ab8da141
      Liwei Lin authored
      ## What changes were proposed in this pull request?
      
      This patch uses `{% highlight lang %}...{% endhighlight %}` to highlight code snippets in the `Structured Streaming Kafka010 integration doc` and the `Spark Streaming Kafka010 integration doc`.
      
      This patch consists of two commits:
      - the first commit fixes only the leading spaces -- this is large
      - the second commit adds the highlight instructions -- this is much simpler and easier to review
      
      ## How was this patch tested?
      
      SKIP_API=1 jekyll build
      
      ## Screenshots
      
      **Before**
      
      ![snip20161101_3](https://cloud.githubusercontent.com/assets/15843379/19894258/47746524-a087-11e6-9a2a-7bff2d428d44.png)
      
      **After**
      
      ![snip20161101_1](https://cloud.githubusercontent.com/assets/15843379/19894324/8bebcd1e-a087-11e6-835b-88c4d2979cfa.png
      
      )
      
      Author: Liwei Lin <lwlin7@gmail.com>
      
      Closes #15715 from lw-lin/doc-highlight-code-snippet.
      
      (cherry picked from commit 98ede494)
      Signed-off-by: default avatarSean Owen <sowen@cloudera.com>
      Unverified
      ab8da141
    • Ryan Blue's avatar
      [SPARK-17532] Add lock debugging info to thread dumps. · 3b624bed
      Ryan Blue authored
      ## What changes were proposed in this pull request?
      
      This adds information to the web UI thread dump page about the JVM locks
      held by threads and the locks that threads are blocked waiting to
      acquire. This should help find cases where lock contention is causing
      Spark applications to run slowly.
      ## How was this patch tested?
      
      Tested by applying this patch and viewing the change in the web UI.
      
      ![thread-lock-info](https://cloud.githubusercontent.com/assets/87915/18493057/6e5da870-79c3-11e6-8c20-f54c18a37544.png
      
      )
      
      Additions:
      - A "Thread Locking" column with the locks held by the thread or that are blocking the thread
      - Links from the a blocked thread to the thread holding the lock
      - Stack frames show where threads are inside `synchronized` blocks, "holding Monitor(...)"
      
      Author: Ryan Blue <blue@apache.org>
      
      Closes #15088 from rdblue/SPARK-17532-add-thread-lock-info.
      
      (cherry picked from commit 2dc04808)
      Signed-off-by: default avatarReynold Xin <rxin@databricks.com>
      3b624bed
    • CodingCat's avatar
      [SPARK-18144][SQL] logging StreamingQueryListener$QueryStartedEvent · 4c4bf87a
      CodingCat authored
      ## What changes were proposed in this pull request?
      
      The PR fixes the bug that the QueryStartedEvent is not logged
      
      the postToAll() in the original code is actually calling StreamingQueryListenerBus.postToAll() which has no listener at all....we shall post by sparkListenerBus.postToAll(s) and this.postToAll() to trigger local listeners as well as the listeners registered in LiveListenerBus
      
      zsxwing
      ## How was this patch tested?
      
      The following snapshot shows that QueryStartedEvent has been logged correctly
      
      ![image](https://cloud.githubusercontent.com/assets/678008/19821553/007a7d28-9d2d-11e6-9f13-49851559cdaa.png
      
      )
      
      Author: CodingCat <zhunansjtu@gmail.com>
      
      Closes #15675 from CodingCat/SPARK-18144.
      
      (cherry picked from commit 85c5424d)
      Signed-off-by: default avatarShixiong Zhu <shixiong@databricks.com>
      4c4bf87a
    • Reynold Xin's avatar
      [SPARK-18192] Support all file formats in structured streaming · 85dd0737
      Reynold Xin authored
      
      ## What changes were proposed in this pull request?
      This patch adds support for all file formats in structured streaming sinks. This is actually a very small change thanks to all the previous refactoring done using the new internal commit protocol API.
      
      ## How was this patch tested?
      Updated FileStreamSinkSuite to add test cases for json, text, and parquet.
      
      Author: Reynold Xin <rxin@databricks.com>
      
      Closes #15711 from rxin/SPARK-18192.
      
      (cherry picked from commit a36653c5)
      Signed-off-by: default avatarReynold Xin <rxin@databricks.com>
      85dd0737
    • Eric Liang's avatar
      [SPARK-18183][SPARK-18184] Fix INSERT [INTO|OVERWRITE] TABLE ... PARTITION for Datasource tables · e6509c24
      Eric Liang authored
      
      There are a couple issues with the current 2.1 behavior when inserting into Datasource tables with partitions managed by Hive.
      
      (1) OVERWRITE TABLE ... PARTITION will actually overwrite the entire table instead of just the specified partition.
      (2) INSERT|OVERWRITE does not work with partitions that have custom locations.
      
      This PR fixes both of these issues for Datasource tables managed by Hive. The behavior for legacy tables or when `manageFilesourcePartitions = false` is unchanged.
      
      There is one other issue in that INSERT OVERWRITE with dynamic partitions will overwrite the entire table instead of just the updated partitions, but this behavior is pretty complicated to implement for Datasource tables. We should address that in a future release.
      
      Unit tests.
      
      Author: Eric Liang <ekl@databricks.com>
      
      Closes #15705 from ericl/sc-4942.
      
      (cherry picked from commit abefe2ec)
      Signed-off-by: default avatarReynold Xin <rxin@databricks.com>
      e6509c24
    • frreiss's avatar
      [SPARK-17475][STREAMING] Delete CRC files if the filesystem doesn't use checksum files · 39d2fdb5
      frreiss authored
      
      ## What changes were proposed in this pull request?
      
      When the metadata logs for various parts of Structured Streaming are stored on non-HDFS filesystems such as NFS or ext4, the HDFSMetadataLog class leaves hidden HDFS-style checksum (CRC) files in the log directory, one file per batch. This PR modifies HDFSMetadataLog so that it detects the use of a filesystem that doesn't use CRC files and removes the CRC files.
      ## How was this patch tested?
      
      Modified an existing test case in HDFSMetadataLogSuite to check whether HDFSMetadataLog correctly removes CRC files on the local POSIX filesystem.  Ran the entire regression suite.
      
      Author: frreiss <frreiss@us.ibm.com>
      
      Closes #15027 from frreiss/fred-17475.
      
      (cherry picked from commit 620da3b4)
      Signed-off-by: default avatarReynold Xin <rxin@databricks.com>
      39d2fdb5
    • Michael Allman's avatar
      [SPARK-17992][SQL] Return all partitions from HiveShim when Hive throws a... · 1bbf9ff6
      Michael Allman authored
      [SPARK-17992][SQL] Return all partitions from HiveShim when Hive throws a metastore exception when attempting to fetch partitions by filter
      
      (Link to Jira issue: https://issues.apache.org/jira/browse/SPARK-17992)
      ## What changes were proposed in this pull request?
      
      We recently added table partition pruning for partitioned Hive tables converted to using `TableFileCatalog`. When the Hive configuration option `hive.metastore.try.direct.sql` is set to `false`, Hive will throw an exception for unsupported filter expressions. For example, attempting to filter on an integer partition column will throw a `org.apache.hadoop.hive.metastore.api.MetaException`.
      
      I discovered this behavior because VideoAmp uses the CDH version of Hive with a Postgresql metastore DB. In this configuration, CDH sets `hive.metastore.try.direct.sql` to `false` by default, and queries that filter on a non-string partition column will fail.
      
      Rather than throw an exception in query planning, this patch catches this exception, logs a warning and returns all table partitions instead. Clients of this method are already expected to handle the possibility that the filters will not be honored.
      ## How was this patch tested?
      
      A unit test was added.
      
      Author: Michael Allman <michael@videoamp.com>
      
      Closes #15673 from mallman/spark-17992-catch_hive_partition_filter_exception.
      1bbf9ff6
    • hyukjinkwon's avatar
      [SPARK-17838][SPARKR] Check named arguments for options and use formatted R... · 1ecfafa0
      hyukjinkwon authored
      [SPARK-17838][SPARKR] Check named arguments for options and use formatted R friendly message from JVM exception message
      
      ## What changes were proposed in this pull request?
      
      This PR proposes to
      - improve the R-friendly error messages rather than raw JVM exception one.
      
        As `read.json`, `read.text`, `read.orc`, `read.parquet` and `read.jdbc` are executed in the same  path with `read.df`, and `write.json`, `write.text`, `write.orc`, `write.parquet` and `write.jdbc` shares the same path with `write.df`, it seems it is safe to call `handledCallJMethod` to handle
        JVM messages.
      -  prevent `zero-length variable name` and prints the ignored options as an warning message.
      
      **Before**
      
      ``` r
      > read.json("path", a = 1, 2, 3, "a")
      Error in env[[name]] <- value :
        zero-length variable name
      ```
      
      ``` r
      > read.json("arbitrary_path")
      Error in invokeJava(isStatic = FALSE, objId$id, methodName, ...) :
        org.apache.spark.sql.AnalysisException: Path does not exist: file:/...;
        at org.apache.spark.sql.execution.datasources.DataSource$$anonfun$12.apply(DataSource.scala:398)
        ...
      
      > read.orc("arbitrary_path")
      Error in invokeJava(isStatic = FALSE, objId$id, methodName, ...) :
        org.apache.spark.sql.AnalysisException: Path does not exist: file:/...;
        at org.apache.spark.sql.execution.datasources.DataSource$$anonfun$12.apply(DataSource.scala:398)
        ...
      
      > read.text("arbitrary_path")
      Error in invokeJava(isStatic = FALSE, objId$id, methodName, ...) :
        org.apache.spark.sql.AnalysisException: Path does not exist: file:/...;
        at org.apache.spark.sql.execution.datasources.DataSource$$anonfun$12.apply(DataSource.scala:398)
        ...
      
      > read.parquet("arbitrary_path")
      Error in invokeJava(isStatic = FALSE, objId$id, methodName, ...) :
        org.apache.spark.sql.AnalysisException: Path does not exist: file:/...;
        at org.apache.spark.sql.execution.datasources.DataSource$$anonfun$12.apply(DataSource.scala:398)
        ...
      ```
      
      ``` r
      > write.json(df, "existing_path")
      Error in invokeJava(isStatic = FALSE, objId$id, methodName, ...) :
        org.apache.spark.sql.AnalysisException: path file:/... already exists.;
        at org.apache.spark.sql.execution.datasources.InsertIntoHadoopFsRelationCommand.run(InsertIntoHadoopFsRelationCommand.scala:68)
      
      > write.orc(df, "existing_path")
      Error in invokeJava(isStatic = FALSE, objId$id, methodName, ...) :
        org.apache.spark.sql.AnalysisException: path file:/... already exists.;
        at org.apache.spark.sql.execution.datasources.InsertIntoHadoopFsRelationCommand.run(InsertIntoHadoopFsRelationCommand.scala:68)
      
      > write.text(df, "existing_path")
      Error in invokeJava(isStatic = FALSE, objId$id, methodName, ...) :
        org.apache.spark.sql.AnalysisException: path file:/... already exists.;
        at org.apache.spark.sql.execution.datasources.InsertIntoHadoopFsRelationCommand.run(InsertIntoHadoopFsRelationCommand.scala:68)
      
      > write.parquet(df, "existing_path")
      Error in invokeJava(isStatic = FALSE, objId$id, methodName, ...) :
        org.apache.spark.sql.AnalysisException: path file:/... already exists.;
        at org.apache.spark.sql.execution.datasources.InsertIntoHadoopFsRelationCommand.run(InsertIntoHadoopFsRelationCommand.scala:68)
      ```
      
      **After**
      
      ``` r
      read.json("arbitrary_path", a = 1, 2, 3, "a")
      Unnamed arguments ignored: 2, 3, a.
      ```
      
      ``` r
      > read.json("arbitrary_path")
      Error in json : analysis error - Path does not exist: file:/...
      
      > read.orc("arbitrary_path")
      Error in orc : analysis error - Path does not exist: file:/...
      
      > read.text("arbitrary_path")
      Error in text : analysis error - Path does not exist: file:/...
      
      > read.parquet("arbitrary_path")
      Error in parquet : analysis error - Path does not exist: file:/...
      ```
      
      ``` r
      > write.json(df, "existing_path")
      Error in json : analysis error - path file:/... already exists.;
      
      > write.orc(df, "existing_path")
      Error in orc : analysis error - path file:/... already exists.;
      
      > write.text(df, "existing_path")
      Error in text : analysis error - path file:/... already exists.;
      
      > write.parquet(df, "existing_path")
      Error in parquet : analysis error - path file:/... already exists.;
      ```
      ## How was this patch tested?
      
      Unit tests in `test_utils.R` and `test_sparkSQL.R`.
      
      Author: hyukjinkwon <gurwls223@gmail.com>
      
      Closes #15608 from HyukjinKwon/SPARK-17838.
      1ecfafa0
  2. Nov 01, 2016
    • Reynold Xin's avatar
      [SPARK-18216][SQL] Make Column.expr public · ad4832a9
      Reynold Xin authored
      ## What changes were proposed in this pull request?
      Column.expr is private[sql], but it's an actually really useful field to have for debugging. We should open it up, similar to how we use QueryExecution.
      
      ## How was this patch tested?
      N/A - this is a simple visibility change.
      
      Author: Reynold Xin <rxin@databricks.com>
      
      Closes #15724 from rxin/SPARK-18216.
      ad4832a9
    • Reynold Xin's avatar
      [SPARK-18025] Use commit protocol API in structured streaming · 77a98162
      Reynold Xin authored
      ## What changes were proposed in this pull request?
      This patch adds a new commit protocol implementation ManifestFileCommitProtocol that follows the existing streaming flow, and uses it in FileStreamSink to consolidate the write path in structured streaming with the batch mode write path.
      
      This deletes a lot of code, and would make it trivial to support other functionalities that are currently available in batch but not in streaming, including all file formats and bucketing.
      
      ## How was this patch tested?
      Should be covered by existing tests.
      
      Author: Reynold Xin <rxin@databricks.com>
      
      Closes #15710 from rxin/SPARK-18025.
      77a98162
    • Joseph K. Bradley's avatar
      [SPARK-18088][ML] Various ChiSqSelector cleanups · 91c33a0c
      Joseph K. Bradley authored
      ## What changes were proposed in this pull request?
      - Renamed kbest to numTopFeatures
      - Renamed alpha to fpr
      - Added missing Since annotations
      - Doc cleanups
      ## How was this patch tested?
      
      Added new standardized unit tests for spark.ml.
      Improved existing unit test coverage a bit.
      
      Author: Joseph K. Bradley <joseph@databricks.com>
      
      Closes #15647 from jkbradley/chisqselector-follow-ups.
      91c33a0c
    • Josh Rosen's avatar
      [SPARK-18182] Expose ReplayListenerBus.read() overload which takes string iterator · b929537b
      Josh Rosen authored
      The `ReplayListenerBus.read()` method is used when implementing a custom `ApplicationHistoryProvider`. The current interface only exposes a `read()` method which takes an `InputStream` and performs stream-to-lines conversion itself, but it would also be useful to expose an overloaded method which accepts an iterator of strings, thereby enabling events to be provided from non-`InputStream` sources.
      
      Author: Josh Rosen <joshrosen@databricks.com>
      
      Closes #15698 from JoshRosen/replay-listener-bus-interface.
      b929537b
    • Josh Rosen's avatar
      [SPARK-17350][SQL] Disable default use of KryoSerializer in Thrift Server · 6e629815
      Josh Rosen authored
      In SPARK-4761 / #3621 (December 2014) we enabled Kryo serialization by default in the Spark Thrift Server. However, I don't think that the original rationale for doing this still holds now that most Spark SQL serialization is now performed via encoders and our UnsafeRow format.
      
      In addition, the use of Kryo as the default serializer can introduce performance problems because the creation of new KryoSerializer instances is expensive and we haven't performed instance-reuse optimizations in several code paths (including DirectTaskResult deserialization).
      
      Given all of this, I propose to revert back to using JavaSerializer as the default serializer in the Thrift Server.
      
      /cc liancheng
      
      Author: Josh Rosen <joshrosen@databricks.com>
      
      Closes #14906 from JoshRosen/disable-kryo-in-thriftserver.
      6e629815
    • hyukjinkwon's avatar
      [SPARK-17764][SQL] Add `to_json` supporting to convert nested struct column to JSON string · 01dd0083
      hyukjinkwon authored
      ## What changes were proposed in this pull request?
      
      This PR proposes to add `to_json` function in contrast with `from_json` in Scala, Java and Python.
      
      It'd be useful if we can convert a same column from/to json. Also, some datasources do not support nested types. If we are forced to save a dataframe into those data sources, we might be able to work around by this function.
      
      The usage is as below:
      
      ``` scala
      val df = Seq(Tuple1(Tuple1(1))).toDF("a")
      df.select(to_json($"a").as("json")).show()
      ```
      
      ``` bash
      +--------+
      |    json|
      +--------+
      |{"_1":1}|
      +--------+
      ```
      ## How was this patch tested?
      
      Unit tests in `JsonFunctionsSuite` and `JsonExpressionsSuite`.
      
      Author: hyukjinkwon <gurwls223@gmail.com>
      
      Closes #15354 from HyukjinKwon/SPARK-17764.
      01dd0083
    • Eric Liang's avatar
      [SPARK-18167] Disable flaky SQLQuerySuite test · cfac17ee
      Eric Liang authored
      We now know it's a persistent environmental issue that is causing this test to sometimes fail. One hypothesis is that some configuration is leaked from another suite, and depending on suite ordering this can cause this test to fail.
      
      I am planning on mining the jenkins logs to try to narrow down which suite could be causing this. For now, disable the test.
      
      Author: Eric Liang <ekl@databricks.com>
      
      Closes #15720 from ericl/disable-flaky-test.
      cfac17ee
    • jiangxingbo's avatar
      [SPARK-18148][SQL] Misleading Error Message for Aggregation Without Window/GroupBy · d0272b43
      jiangxingbo authored
      ## What changes were proposed in this pull request?
      
      Aggregation Without Window/GroupBy expressions will fail in `checkAnalysis`, the error message is a bit misleading, we should generate a more specific error message for this case.
      
      For example,
      
      ```
      spark.read.load("/some-data")
        .withColumn("date_dt", to_date($"date"))
        .withColumn("year", year($"date_dt"))
        .withColumn("week", weekofyear($"date_dt"))
        .withColumn("user_count", count($"userId"))
        .withColumn("daily_max_in_week", max($"user_count").over(weeklyWindow))
      )
      ```
      
      creates the following output:
      
      ```
      org.apache.spark.sql.AnalysisException: expression '`randomColumn`' is neither present in the group by, nor is it an aggregate function. Add to group by or wrap in first() (or first_value) if you don't care which value you get.;
      ```
      
      In the error message above, `randomColumn` doesn't appear in the query(acturally it's added by function `withColumn`), so the message is not enough for the user to address the problem.
      ## How was this patch tested?
      
      Manually test
      
      Before:
      
      ```
      scala> spark.sql("select col, count(col) from tbl")
      org.apache.spark.sql.AnalysisException: expression 'tbl.`col`' is neither present in the group by, nor is it an aggregate function. Add to group by or wrap in first() (or first_value) if you don't care which value you get.;;
      ```
      
      After:
      
      ```
      scala> spark.sql("select col, count(col) from tbl")
      org.apache.spark.sql.AnalysisException: grouping expressions sequence is empty, and 'tbl.`col`' is not an aggregate function. Wrap '(count(col#231L) AS count(col)#239L)' in windowing function(s) or wrap 'tbl.`col`' in first() (or first_value) if you don't care which value you get.;;
      ```
      
      Also add new test sqls in `group-by.sql`.
      
      Author: jiangxingbo <jiangxb1987@gmail.com>
      
      Closes #15672 from jiangxb1987/groupBy-empty.
      d0272b43
    • Ergin Seyfe's avatar
      [SPARK-18189][SQL] Fix serialization issue in KeyValueGroupedDataset · 8a538c97
      Ergin Seyfe authored
      ## What changes were proposed in this pull request?
      Likewise [DataSet.scala](https://github.com/apache/spark/blob/master/sql/core/src/main/scala/org/apache/spark/sql/Dataset.scala#L156) KeyValueGroupedDataset should mark the queryExecution as transient.
      
      As mentioned in the Jira ticket, without transient we saw serialization issues like
      
      ```
      Caused by: java.io.NotSerializableException: org.apache.spark.sql.execution.QueryExecution
      Serialization stack:
              - object not serializable (class: org.apache.spark.sql.execution.QueryExecution, value: ==
      ```
      
      ## How was this patch tested?
      
      Run the query which is specified in the Jira ticket before and after:
      ```
      val a = spark.createDataFrame(sc.parallelize(Seq((1,2),(3,4)))).as[(Int,Int)]
      val grouped = a.groupByKey(
      {x:(Int,Int)=>x._1}
      )
      val mappedGroups = grouped.mapGroups((k,x)=>
      {(k,1)}
      )
      val yyy = sc.broadcast(1)
      val last = mappedGroups.rdd.map(xx=>
      { val simpley = yyy.value 1 }
      )
      ```
      
      Author: Ergin Seyfe <eseyfe@fb.com>
      
      Closes #15706 from seyfe/keyvaluegrouped_serialization.
      8a538c97
    • Liwei Lin's avatar
      [SPARK-18103][FOLLOW-UP][SQL][MINOR] Rename `MetadataLogFileCatalog` to `MetadataLogFileIndex` · 8cdf143f
      Liwei Lin authored
      ## What changes were proposed in this pull request?
      
      This is a follow-up to https://github.com/apache/spark/pull/15634.
      
      ## How was this patch tested?
      
      N/A
      
      Author: Liwei Lin <lwlin7@gmail.com>
      
      Closes #15712 from lw-lin/18103.
      8cdf143f
    • Zheng RuiFeng's avatar
      [SPARK-17848][ML] Move LabelCol datatype cast into Predictor.fit · 8ac09108
      Zheng RuiFeng authored
      ## What changes were proposed in this pull request?
      
      1, move cast to `Predictor`
      2, and then, remove unnecessary cast
      ## How was this patch tested?
      
      existing tests
      
      Author: Zheng RuiFeng <ruifengz@foxmail.com>
      
      Closes #15414 from zhengruifeng/move_cast.
      8ac09108
    • Herman van Hovell's avatar
      0cba535a
    • eyal farago's avatar
      [SPARK-16839][SQL] redundant aliases after cleanupAliases · 5441a626
      eyal farago authored
      ## What changes were proposed in this pull request?
      
      Simplify struct creation, especially the aspect of `CleanupAliases` which missed some aliases when handling trees created by `CreateStruct`.
      
      This PR includes:
      
      1. A failing test (create struct with nested aliases, some of the aliases survive `CleanupAliases`).
      2. A fix that transforms `CreateStruct` into a `CreateNamedStruct` constructor, effectively eliminating `CreateStruct` from all expression trees.
      3. A `NamePlaceHolder` used by `CreateStruct` when column names cannot be extracted from unresolved `NamedExpression`.
      4. A new Analyzer rule that resolves `NamePlaceHolder` into a string literal once the `NamedExpression` is resolved.
      5. `CleanupAliases` code was simplified as it no longer has to deal with `CreateStruct`'s top level columns.
      
      ## How was this patch tested?
      
      running all tests-suits in package org.apache.spark.sql, especially including the analysis suite, making sure added test initially fails, after applying suggested fix rerun the entire analysis package successfully.
      
      modified few tests that expected `CreateStruct` which is now transformed into `CreateNamedStruct`.
      
      Credit goes to hvanhovell for assisting with this PR.
      
      Author: eyal farago <eyal farago>
      Author: eyal farago <eyal.farago@gmail.com>
      Author: Herman van Hovell <hvanhovell@databricks.com>
      Author: Eyal Farago <eyal.farago@actimize.com>
      Author: Hyukjin Kwon <gurwls223@gmail.com>
      Author: eyalfa <eyal.farago@gmail.com>
      
      Closes #14444 from eyalfa/SPARK-16839_redundant_aliases_after_cleanupAliases.
      5441a626
    • Herman van Hovell's avatar
      [SPARK-17996][SQL] Fix unqualified catalog.getFunction(...) · f7c145d8
      Herman van Hovell authored
      ## What changes were proposed in this pull request?
      
      Currently an unqualified `getFunction(..)`call returns a wrong result; the returned function is shown as temporary function without a database. For example:
      
      ```
      scala> sql("create function fn1 as 'org.apache.hadoop.hive.ql.udf.generic.GenericUDFAbs'")
      res0: org.apache.spark.sql.DataFrame = []
      
      scala> spark.catalog.getFunction("fn1")
      res1: org.apache.spark.sql.catalog.Function = Function[name='fn1', className='org.apache.hadoop.hive.ql.udf.generic.GenericUDFAbs', isTemporary='true']
      ```
      
      This PR fixes this by adding database information to ExpressionInfo (which is used to store the function information).
      ## How was this patch tested?
      
      Added more thorough tests to `CatalogSuite`.
      
      Author: Herman van Hovell <hvanhovell@databricks.com>
      
      Closes #15542 from hvanhovell/SPARK-17996.
      f7c145d8
    • Wang Lei's avatar
      [SPARK-18114][MESOS] Fix mesos cluster scheduler generage command option error · 9b377aa4
      Wang Lei authored
      ## What changes were proposed in this pull request?
      
      Enclose --conf option value with "" to support multi value configs like spark.driver.extraJavaOptions, without "", driver will fail to start.
      ## How was this patch tested?
      
      Jenkins Tests.
      
      Test in our production environment, also unit tests, It is a very small change.
      
      Author: Wang Lei <lei.wang@kongming-inc.com>
      
      Closes #15643 from LeightonWong/messos-cluster.
      Unverified
      9b377aa4
    • Sandeep Singh's avatar
      [SPARK-16881][MESOS] Migrate Mesos configs to use ConfigEntry · ec6f479b
      Sandeep Singh authored
      ## What changes were proposed in this pull request?
      
      Migrate Mesos configs to use ConfigEntry
      ## How was this patch tested?
      
      Jenkins Tests
      
      Author: Sandeep Singh <sandeep@techaddict.me>
      
      Closes #15654 from techaddict/SPARK-16881.
      Unverified
      ec6f479b
    • Charles Allen's avatar
      [SPARK-15994][MESOS] Allow enabling Mesos fetch cache in coarse executor backend · e34b4e12
      Charles Allen authored
      Mesos 0.23.0 introduces a Fetch Cache feature http://mesos.apache.org/documentation/latest/fetcher/ which allows caching of resources specified in command URIs.
      
      This patch:
      - Updates the Mesos shaded protobuf dependency to 0.23.0
      - Allows setting `spark.mesos.fetcherCache.enable` to enable the fetch cache for all specified URIs. (URIs must be specified for the setting to have any affect)
      - Updates documentation for Mesos configuration with the new setting.
      
      This patch does NOT:
      - Allow for per-URI caching configuration. The cache setting is global to ALL URIs for the command.
      
      Author: Charles Allen <charles@allen-net.com>
      
      Closes #13713 from drcrallen/SPARK15994.
      Unverified
      e34b4e12
    • wangzhenhua's avatar
      [SPARK-18111][SQL] Wrong ApproximatePercentile answer when multiple records have the minimum value · cb80edc2
      wangzhenhua authored
      ## What changes were proposed in this pull request?
      
      When multiple records have the minimum value, the answer of ApproximatePercentile is wrong.
      ## How was this patch tested?
      
      add a test case
      
      Author: wangzhenhua <wangzhenhua@huawei.com>
      
      Closes #15641 from wzhfy/percentile.
      Unverified
      cb80edc2
    • Dongjoon Hyun's avatar
      [MINOR][DOC] Remove spaces following slashs · 623fc7fc
      Dongjoon Hyun authored
      ## What changes were proposed in this pull request?
      
      This PR merges multiple lines enumerating items in order to remove the redundant spaces following slashes in [Structured Streaming Programming Guide in 2.0.2-rc1](http://people.apache.org/~pwendell/spark-releases/spark-2.0.2-rc1-docs/structured-streaming-programming-guide.html).
      - Before: `Scala/ Java/ Python`
      - After: `Scala/Java/Python`
      ## How was this patch tested?
      
      Manual by the followings because this is documentation update.
      
      ```
      cd docs
      SKIP_API=1 jekyll build
      ```
      
      Author: Dongjoon Hyun <dongjoon@apache.org>
      
      Closes #15686 from dongjoon-hyun/minor_doc_space.
      Unverified
      623fc7fc
    • Liang-Chi Hsieh's avatar
      [SPARK-18107][SQL] Insert overwrite statement runs much slower in spark-sql... · dd85eb54
      Liang-Chi Hsieh authored
      [SPARK-18107][SQL] Insert overwrite statement runs much slower in spark-sql than it does in hive-client
      
      ## What changes were proposed in this pull request?
      
      As reported on the jira, insert overwrite statement runs much slower in Spark, compared with hive-client.
      
      It seems there is a patch [HIVE-11940](https://github.com/apache/hive/commit/ba21806b77287e237e1aa68fa169d2a81e07346d) which largely improves insert overwrite performance on Hive. HIVE-11940 is patched after Hive 2.0.0.
      
      Because Spark SQL uses older Hive library, we can not benefit from such improvement.
      
      The reporter verified that there is also a big performance gap between Hive 1.2.1 (520.037 secs) and Hive 2.0.1 (35.975 secs) on insert overwrite execution.
      
      Instead of upgrading to Hive 2.0 in Spark SQL, which might not be a trivial task, this patch provides an approach to delete the partition before asking Hive to load data files into the partition.
      
      Note: The case reported on the jira is insert overwrite to partition. Since `Hive.loadTable` also uses the function to replace files, insert overwrite to table should has the same issue. We can take the same approach to delete the table first. I will upgrade this to include this.
      ## How was this patch tested?
      
      Jenkins tests.
      
      There are existing tests using insert overwrite statement. Those tests should be passed. I added a new test to specially test insert overwrite into partition.
      
      For performance issue, as I don't have Hive 2.0 environment, this needs the reporter to verify it. Please refer to the jira.
      
      Please review https://cwiki.apache.org/confluence/display/SPARK/Contributing+to+Spark before opening a pull request.
      
      Author: Liang-Chi Hsieh <viirya@gmail.com>
      
      Closes #15667 from viirya/improve-hive-insertoverwrite.
      dd85eb54
    • Reynold Xin's avatar
      [SPARK-18024][SQL] Introduce an internal commit protocol API · d9d14650
      Reynold Xin authored
      ## What changes were proposed in this pull request?
      This patch introduces an internal commit protocol API that is used by the batch data source to do write commits. It currently has only one implementation that uses Hadoop MapReduce's OutputCommitter API. In the future, this commit API can be used to unify streaming and batch commits.
      
      ## How was this patch tested?
      Should be covered by existing write tests.
      
      Author: Reynold Xin <rxin@databricks.com>
      Author: Eric Liang <ekl@databricks.com>
      
      Closes #15707 from rxin/SPARK-18024-2.
      d9d14650
  3. Oct 31, 2016
    • Eric Liang's avatar
      [SPARK-18167][SQL] Retry when the SQLQuerySuite test flakes · 7d6c8715
      Eric Liang authored
      ## What changes were proposed in this pull request?
      
      This will re-run the flaky test a few times after it fails. This will help determine if it's due to nondeterministic test setup, or because of some environment issue (e.g. leaked config from another test).
      
      cc yhuai
      
      Author: Eric Liang <ekl@databricks.com>
      
      Closes #15708 from ericl/spark-18167-3.
      7d6c8715
    • Eric Liang's avatar
      [SPARK-18087][SQL] Optimize insert to not require REPAIR TABLE · efc254a8
      Eric Liang authored
      ## What changes were proposed in this pull request?
      
      When inserting into datasource tables with partitions managed by the hive metastore, we need to notify the metastore of newly added partitions. Previously this was implemented via `msck repair table`, but this is more expensive than needed.
      
      This optimizes the insertion path to add only the updated partitions.
      ## How was this patch tested?
      
      Existing tests (I verified manually that tests fail if the repair operation is omitted).
      
      Author: Eric Liang <ekl@databricks.com>
      
      Closes #15633 from ericl/spark-18087.
      efc254a8
    • Eric Liang's avatar
      [SPARK-18167][SQL] Also log all partitions when the SQLQuerySuite test flakes · 6633b97b
      Eric Liang authored
      ## What changes were proposed in this pull request?
      
      One possibility for this test flaking is that we have corrupted the partition schema somehow in the tests, which causes the cast to decimal to fail in the call. This should at least show us the actual partition values.
      
      ## How was this patch tested?
      
      Run it locally, it prints out something like `ArrayBuffer(test(partcol=0), test(partcol=1), test(partcol=2), test(partcol=3), test(partcol=4))`.
      
      Author: Eric Liang <ekl@databricks.com>
      
      Closes #15701 from ericl/print-more-info.
      6633b97b
    • Shixiong Zhu's avatar
      [SPARK-18030][TESTS] Fix flaky FileStreamSourceSuite by not deleting the files · de3f87fa
      Shixiong Zhu authored
      ## What changes were proposed in this pull request?
      
      The test `when schema inference is turned on, should read partition data` should not delete files because the source maybe is listing files. This PR just removes the delete actions since they are not necessary.
      
      ## How was this patch tested?
      
      Jenkins
      
      Author: Shixiong Zhu <shixiong@databricks.com>
      
      Closes #15699 from zsxwing/SPARK-18030.
      de3f87fa
Loading