小编aux*_*xdx的帖子

如何检查spark数据帧是否为空

现在,我必须用来df.count > 0检查它是否DataFrame为空.但它效率低下.有没有更好的方法来做到这一点.

谢谢.

PS:我想检查它是否为空,以便我只保存,DataFrame如果它不是空的

apache-spark apache-spark-sql

82
推荐指数
9
解决办法
6万
查看次数

创建 AWS 应用程序负载均衡器规则,无需使用额外的反向代理(如 nginx、httpd)即可修剪请求的前缀

基本上,我有几个服务。我想将带有前缀“/secured”的每个请求转发到 server1 端口 80,并将所有其他请求转发到服务器 2 端口 80。问题是在 server1 上,我正在运行接受没有“/secured”前缀的请求的服务。换句话说,我希望每次请求,例如“转发http://example.com/secured/api/getUser ”到server1为“ http://example.com/api/getUser ”(删除/从请求固定”小路)。

使用 AWS ALB,当前请求以http://example.com/secured/api/getUser 的形式发送;这迫使我更新我的 server1 的代码,以便代码处理带有 /secured 前缀的请求,这看起来不太好。

有没有什么简单的方法可以用 ALB 解决这个问题?

谢谢。

reverse-proxy amazon-web-services elastic-load-balancer

10
推荐指数
2
解决办法
3953
查看次数

使用Dispatcher的Spark Mesos集群模式

我只有一台机器,并希望使用mesos集群模式运行spark作业.使用节点集群运行可能更有意义,但我主要想首先测试mesos以检查它是否能够更有效地利用资源(同时运行多个spark作业而不进行静态分区).我尝试了很多方法但没有成功.这是我做的:

  1. 构建mesos并运行mesos主站和从站(同一台机器中的2个从站).

    sudo ./bin/mesos-master.sh --ip=127.0.0.1 --work_dir=/var/lib/mesos
    sudo ./bin/mesos-slave.sh --master=127.0.0.1:5050 --port=5051 --work_dir=/tmp/mesos1
    sudo ./bin/mesos-slave.sh --master=127.0.0.1:5050 --port=5052 --work_dir=/tmp/mesos2
    
    Run Code Online (Sandbox Code Playgroud)
  2. 运行spark-mesos-dispatcher

    sudo ./sbin/start-mesos-dispatcher.sh --master mesos://localhost:5050
    
    Run Code Online (Sandbox Code Playgroud)
  3. 使用调度程序作为主URL提交应用程序.

    spark-submit  --master mesos://localhost:7077 <other-config> <jar file>
    
    Run Code Online (Sandbox Code Playgroud)

但它不起作用:

    E0925 17:30:30.158846 807608320 socket.hpp:174] Shutdown failed on fd=61: Socket is not connected [57]
    E0925 17:30:30.159545 807608320 socket.hpp:174] Shutdown failed on fd=62: Socket is not connected [57]
Run Code Online (Sandbox Code Playgroud)

如果我使用spark-submit --deploy-mode集群,那么我收到另一条错误消息:

    Exception in thread "main" org.apache.spark.deploy.rest.SubmitRestConnectionException: Unable to connect to server
Run Code Online (Sandbox Code Playgroud)

如果我不使用调度程序但直接使用mesos master url,它可以正常工作: - master mesos:// localhost:5050(客户端模式).根据文档,Mesos集群不支持集群模式,但它们在此处为集群模式提供了另一条指令.这有点令人困惑?我的问题是:

  1. 我怎么能搞定它?
  2. 如果我直接从主节点提交app/jar,我应该使用客户端模式而不是群集模式吗? …

mesos apache-spark

6
推荐指数
1
解决办法
2912
查看次数

将 Spark 数据帧保存为 Google Cloud Storage 中的 parquet 文件

我正在尝试将 Spark 数据帧保存到 Google Cloud Storage。我们可以将 parquet 格式的数据帧保存到 S3,但由于我们的服务器是 Google Compute Engine,因此到 S3 的数据传输成本会很高。我想知道谷歌云存储是否可以有类似的功能?以下是我在 S3 中所做的操作:

添加依赖到build.sbt:

"net.java.dev.jets3t" % "jets3t" % "0.9.4",
"com.amazonaws" % "aws-java-sdk" % "1.10.16"
Run Code Online (Sandbox Code Playgroud)

在主代码中使用它:

val sc = new SparkContext(sparkConf)
sc.hadoopConfiguration.set("fs.s3a.awsAccessKeyId", conf.getString("s3.awsAccessKeyId"))
sc.hadoopConfiguration.set("fs.s3a.awsSecretAccessKey", conf.getString("s3.awsSecretAccessKey"))

val df = sqlContext.read.parquet("s3a://.../*") //read file
df.write.mode(SaveMode.Append).parquet(s3FileName) //write file
Run Code Online (Sandbox Code Playgroud)

最后,将其与 Spark-Submit 一起使用

spark-submit --conf spark.hadoop.fs.s3a.impl=org.apache.hadoop.fs.s3native.NativeS3FileSystem 
--conf spark.hadoop.fs.s3.impl=org.apache.hadoop.fs.s3.S3FileSystem
Run Code Online (Sandbox Code Playgroud)

我试图在互联网上寻找类似的指南,但似乎没有?任何人都可以建议我如何完成它吗?

谢谢。

scala google-cloud-storage apache-spark parquet apache-spark-sql

5
推荐指数
1
解决办法
3630
查看次数

Spark Executor:检测到托管内存泄漏

我正在使用mesos集群来部署spark job(客户端模式).我有三台服务器,能够运行spark工作.但是,过了一会儿(几天),我收到了错误:

5/11/03 19:55:50 ERROR Executor: Managed memory leak detected; size = 33554432 bytes, TID = 387939
15/11/03 19:55:50 ERROR Executor: Exception in task 2.1 in stage 6534.0 (TID 387939)
java.io.FileNotFoundException: /tmp/blockmgr-3acec504-4a55-4aa8-a3e5-dda97ce5d055/03/temp_shuffle_cb37f147-c055-4014-a6ae-fd505cb49f57 (Too many open files)
    at java.io.FileOutputStream.open(Native Method)
    at java.io.FileOutputStream.<init>(FileOutputStream.java:221)
    at org.apache.spark.storage.DiskBlockObjectWriter.open(DiskBlockObjectWriter.scala:88)
    at org.apache.spark.shuffle.sort.BypassMergeSortShuffleWriter.insertAll(BypassMergeSortShuffleWriter.java:110)
    at org.apache.spark.shuffle.sort.SortShuffleWriter.write(SortShuffleWriter.scala:73)
    at org.apache.spark.scheduler.ShuffleMapTask.runTask(ShuffleMapTask.scala:73)
    at org.apache.spark.scheduler.ShuffleMapTask.runTask(ShuffleMapTask.scala:41)
    at org.apache.spark.scheduler.Task.run(Task.scala:88)
    at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:214)
    at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1145)
    at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:615)
    at java.lang.Thread.run(Thread.java:745)
Run Code Online (Sandbox Code Playgroud)

现在,它导致所有流批处理排队并显示为"处理"(4042 /流/).在手动重新启动spark作业并再次重新提交之前,他们都无法继续.

我的火花工作只是从kafka读取数据并对mongo进行一些更新(有很多更新查询通过;但我将火花流持续时间配置为大约5分钟;所以它不应该导致问题).

过了一会儿,因为没有工作能够成功; spark-kafka阅读器开始显示错误:

ERROR Executor: Exception in task 5.3 in stage 7561.0 (TID 392220)
org.apache.spark.SparkException: Couldn't connect …
Run Code Online (Sandbox Code Playgroud)

memory-leaks apache-kafka mesos apache-spark spark-streaming

5
推荐指数
0
解决办法
5189
查看次数

Max Double Slice Sum codility O(1)空间复杂性失败性能测试用例

我试图弄清楚为什么下面的解决方案在代码网站中针对"Max Double Slice Sum"问题的单个性能测试案例失败了:https://codility.com/demo/take-sample-test/max_double_slice_sum

还有另一种解决方案O(n)空间复杂度更容易理解:最大双切片和.但我只是想知道为什么这个O(1)解决方案不起作用.以下是实际代码:

import java.util.*;

class Solution {
    public int solution(int[] A) {      
        long maxDS = 0;
        long maxDSE = 0;
        long maxS = A[1];

        for(int i=2; i<A.length-1; ++i){
            //end at i-index
            maxDSE = Math.max(maxDSE+A[i], maxS);            
            maxDS = Math.max(maxDS, maxDSE);            
            maxS = Math.max(A[i], maxS + A[i]);            
        }

        return (int)maxDS;
    }
}
Run Code Online (Sandbox Code Playgroud)

这个想法很简单如下:

  • 该问题可以作为查找max(A [i] + A [i + 1] + ... + A [j] -A [m]); 1 <= I <= M <= j的<= N-2; …

java algorithm max

3
推荐指数
1
解决办法
1564
查看次数

处理火花流中的数据库连接

我不确定我是否正确理解火花处理数据库连接如何以及如何可靠地使用大量数据库更新操作内部火花而不会搞砸火花作业.这是我一直在使用的代码片段(为了便于说明):

val driver = new MongoDriver
val hostList: List[String] = conf.getString("mongo.hosts").split(",").toList
val connection = driver.connection(hostList)
val mongodb = connection(conf.getString("mongo.db"))
val dailyInventoryCol = mongodb[BSONCollection](conf.getString("mongo.collections.dailyInventory"))

val stream: InputDStream[(String,String)] = KafkaUtils.createDirectStream[String, String, StringDecoder, StringDecoder, (String, String)](
  ssc, kafkaParams, fromOffsets,
  (mmd: MessageAndMetadata[String, String]) => (mmd.topic, mmd.message()));

def processRDD(rddElem: RDD[(String, String)]): Unit = {
    val df = rdd.map(line => {
        ...
    }).flatMap(x => x).toDF()

    if (!isEmptyDF(df)) {
        var mongoF: Seq[Future[dailyInventoryCol.BatchCommands.FindAndModifyCommand.FindAndModifyResult]] = Seq();

        val dfF2 = df.groupBy($"CountryCode", $"Width", $"Height", $"RequestType", $"Timestamp").agg(sum($"Frequency")).collect().map(row => {
        val countryCode = …
Run Code Online (Sandbox Code Playgroud)

mesos apache-spark spark-streaming spark-dataframe

3
推荐指数
1
解决办法
3353
查看次数