script_path=$(cd `dirname $0`; pwd)
cd ${script_path}
显示的就是正在执行的脚本所在的文件夹的绝对路径
2017年11月8日星期三
2017年11月2日星期四
pyspark sql使用udf后yarn模式运行卡住
https://stackoverflow.com/questions/35157322/spark-dataframe-in-python-execution-stuck-when-using-udfs
我们用的pyspark 2.1.0版本,udf还有各种各样的问题,而且性能很差,只能转成rdd再做操作?
我们用的pyspark 2.1.0版本,udf还有各种各样的问题,而且性能很差,只能转成rdd再做操作?
pyspark sql user defined function
参考https://docs.databricks.com/spark/latest/spark-sql/udf-in-python.html
举例:
>>> a = [{'a': 'a', 'b': 1}, {'a': 'aa', 'b': 2}]
>>> df = spark.createDataFrame(a)
[Row(a=u'a', b=1), Row(a=u'aa', b=2)]
>>> def func(str):
... return len(str) > 1
...
>>> from pyspark.sql.functions import udf
>>> from pyspark.sql.types import BooleanType
>>> func_udf = udf(func, BooleanType())
>>> df2 = df.filter(func_udf(df['a']))
>>> df2.collect()
[Row(a=u'aa', b=2)]
举例:
>>> a = [{'a': 'a', 'b': 1}, {'a': 'aa', 'b': 2}]
>>> df = spark.createDataFrame(a)
[Row(a=u'a', b=1), Row(a=u'aa', b=2)]
>>> def func(str):
... return len(str) > 1
...
>>> from pyspark.sql.functions import udf
>>> from pyspark.sql.types import BooleanType
>>> func_udf = udf(func, BooleanType())
>>> df2 = df.filter(func_udf(df['a']))
>>> df2.collect()
[Row(a=u'aa', b=2)]
2017年11月1日星期三
解决github clone速度特别慢
参考http://www.jianshu.com/p/5e74b1042b70
vim ~/.gitconfig,添加:
[http]
proxy = socks5://127.0.0.1:8080
[https]
proxy = socks5://127.0.0.1:8080
使用 ssh -D 127.0.0.1:8080 username@服务器名 命令开启sock5端口转发。
vim ~/.gitconfig,添加:
[http]
proxy = socks5://127.0.0.1:8080
[https]
proxy = socks5://127.0.0.1:8080
使用 ssh -D 127.0.0.1:8080 username@服务器名 命令开启sock5端口转发。
2017年10月31日星期二
spark合并小文件
spark使用FileUtil.copyMerge来进行小文件合并:https://hadoop.apache.org/docs/r2.7.1/api/org/apache/hadoop/fs/FileUtil.html
pyspark中dataframe union的一个问题
>>> a = [{'a': 1, 'b': 2}]
>>> x = spark.createDataFrame(a)
>>> b = sc.parallelize([(3, 4)])
>>> y = spark.createDataFrame(b, ['b', 'a'])
>>> x.collect()
[Row(a=1, b=2)]
>>> y.collect()
[Row(b=3, a=4)]
>>> z = x.union(y)
>>> z.collect()
[Row(a=1, b=2), Row(a=3, b=4)]
正确的结果应该是[Row(a=1, b=2), Row(a=4, b=3)],但实际输出的结果第二个Row的a和b反了。猜测DataFrame的union是按照顺序来的,并不是按照column的名称对应的。
Also as standard in SQL, this function resolves columns by position (not by name). Spark 2.3提供了unionByName可以解决问题,目前解决办法是把x与y的字段名排序要一样才行。
>>> x = spark.createDataFrame(a)
>>> b = sc.parallelize([(3, 4)])
>>> y = spark.createDataFrame(b, ['b', 'a'])
>>> x.collect()
[Row(a=1, b=2)]
>>> y.collect()
[Row(b=3, a=4)]
>>> z = x.union(y)
>>> z.collect()
[Row(a=1, b=2), Row(a=3, b=4)]
正确的结果应该是[Row(a=1, b=2), Row(a=4, b=3)],但实际输出的结果第二个Row的a和b反了。猜测DataFrame的union是按照顺序来的,并不是按照column的名称对应的。
Also as standard in SQL, this function resolves columns by position (not by name). Spark 2.3提供了unionByName可以解决问题,目前解决办法是把x与y的字段名排序要一样才行。
2017年10月30日星期一
spark读写mongodb的一个问题
从mongodb的某个collection中读取了df,做了一些操作后又overwrite写回该collection会有问题。因为在写的时候才action,猜测可能因为分布式的同时读写造成的问题。
问题确认:
将df cache,在回写之前先做一次action,让结果缓存到内存,然后再写mongo没有问题。
解决:
从一个collection读,写到另一个 collection
问题确认:
将df cache,在回写之前先做一次action,让结果缓存到内存,然后再写mongo没有问题。
解决:
从一个collection读,写到另一个 collection
订阅:
博文 (Atom)