spark日志+hivesql
windows本地读取hive,需要在resource里面将集群中的hive-site.xml下载下来。
<?xml version="1.0" encoding="UTF-8"?>
hive.metastore.local
false
hive.metastore.uris
thrift://bn00:9083
hive.metastore.client.socket.timeout
300
hive.metastore.warehouse.dir
/user/hive/warehouse
hive.warehouse.subdir.inherit.perms
true
hive.auto.convert.join
true
hive.auto.convert.join.noconditionaltask.size
20971520
hive.optimize.bucketmapjoin.sortedmerge
false
hive.smbjoin.cache.rows
10000
hive.server2.logging.operation.enabled
true
hive.server2.logging.operation.log.location
/var/log/hive/operation_logs
mapred.reduce.tasks
-1
hive.exec.reducers.bytes.per.reducer
67108864
hive.exec.copyfile.maxsize
33554432
hive.exec.reducers.max
1099
hive.vectorized.groupby.checkinterval
4096
hive.vectorized.groupby.flush.percent
0.1
hive.compute.query.using.stats
true
hive.vectorized.execution.enabled
true
hive.vectorized.execution.reduce.enabled
false
hive.merge.mapfiles
true
hive.merge.mapredfiles
false
hive.cbo.enable
true
hive.fetch.task.conversion
minimal
hive.fetch.task.conversion.threshold
268435456
hive.limit.pushdown.memory.usage
0.1
hive.merge.sparkfiles
true
hive.merge.smallfiles.avgsize
16777216
hive.merge.size.per.task
268435456
hive.optimize.reducededuplication
true
hive.optimize.reducededuplication.min.reducer
4
hive.map.aggr
true
hive.map.aggr.hash.percentmemory
0.5
hive.optimize.sort.dynamic.partition
false
hive.execution.engine
mr
spark.executor.memory
1277794713
spark.driver.memory
966367641
spark.executor.cores
6
spark.yarn.driver.memoryOverhead
102
spark.yarn.executor.memoryOverhead
135
spark.dynamicAllocation.enabled
true
spark.dynamicAllocation.initialExecutors
1
spark.dynamicAllocation.minExecutors
1
spark.dynamicAllocation.maxExecutors
2147483647
hive.metastore.execute.setugi
true
hive.support.concurrency
true
hive.zookeeper.quorum
bn00,bn01,bn02
hive.zookeeper.client.port
2181
hive.zookeeper.namespace
hive_zookeeper_namespace_hive
hbase.zookeeper.quorum
bn00,bn01,bn02
hbase.zookeeper.property.clientPort
2181
hive.cluster.delegation.token.store.class
org.apache.hadoop.hive.thrift.MemoryTokenStore
hive.server2.enable.doAs
true
hive.server2.use.SSL
false
spark.shuffle.service.enabled
true
代码部分如下:
import java.util.ArrayList;
import java.util.List;
import org.apache.log4j.Level;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
import org.apache.spark.SparkConf;
import org.apache.spark.api.java.JavaRDD;
import org.apache.spark.api.java.function.FlatMapFunction;
import org.apache.spark.api.java.function.Function;
import org.apache.spark.api.java.function.Function2;
import org.apache.spark.api.java.function.PairFunction;
import org.apache.spark.api.java.function.VoidFunction;
import org.apache.spark.sql.DataFrame;
import org.apache.spark.sql.Row;
import org.apache.spark.sql.RowFactory;
import org.apache.spark.sql.SQLContext;
import org.apache.spark.sql.hive.HiveContext;
import org.apache.spark.sql.types.DataTypes;
import org.apache.spark.sql.types.StructField;
import org.apache.spark.sql.types.StructType;
import org.apache.spark.streaming.Durations;
import org.apache.spark.streaming.api.java.JavaDStream;
import org.apache.spark.streaming.api.java.JavaStreamingContext;
import scala.Tuple2;
import com.google.common.collect.Lists;
public class HiveAndSparkSQLApp {
private static final Logger logger = LogManager.getLogger(App.class);
static {
// 设置日志级别清理
org.apache.log4j.Logger.getLogger("org.apache.spark").setLevel(Level.WARN);
org.apache.log4j.Logger.getLogger("org.eclipse.jetty.server").setLevel(Level.OFF);
}
@SuppressWarnings("serial")
public static void main(String[] args) {
// 调试环境,spark UI:http://localhost:4040/executors/
SparkConf conf = new SparkConf().setMaster("local[4]").setAppName("test")
.set("spark.testing.memory", "1147480000");
// spark streaming context
JavaStreamingContext jssc = new JavaStreamingContext(conf, Durations.seconds(5));
// spark hive context
final HiveContext hiveContext = new HiveContext(jssc.sparkContext());
// spark SQL context
// final SQLContext sqlContext = SQLContext.getOrCreate(jssc.sparkContext().sc());
/**
* 远程的socket监听
* 在节点上,执行nc -lk 9998
* 若节点上没有安装nc工具,执行yum install nc.x86_64
* 之后直接发送消息即可
*/
JavaDStream lines = jssc.socketTextStream("node0", 9998);
lines.foreachRDD(new VoidFunction>() {
@Override
public void call(JavaRDD rdd) throws Exception {
// SQLContext sqlContext = SQLContext.getOrCreate(rdd.context());
JavaRDD rowRDD = rdd.map(new Function() {
@Override
public Row call(String t) throws Exception {
String[] splited = new String[] { System.currentTimeMillis() + "",
System.currentTimeMillis() + "", System.currentTimeMillis() + "" };
// 1.Row构建
return RowFactory.create(Long.valueOf(splited[0]), splited[1], Long.valueOf(splited[2]));
}
});
// 2.DF metadata专用结构体
// 对Row具体指定元数据信息。
List structFields = new ArrayList();
// 列名称 列的具体类型(Integer Or String) 是否为空一般为true,实际在开发环境是通过for循环,而不是手动添加
structFields.add(DataTypes.createStructField("id", DataTypes.LongType, true));
structFields.add(DataTypes.createStructField("name", DataTypes.StringType, true));
structFields.add(DataTypes.createStructField("age", DataTypes.LongType, true));
// 构建StructType,用于最后DataFrame元数据的描述
StructType structType = DataTypes.createStructType(structFields);
// 3.构建DF
DataFrame personsDF = hiveContext.createDataFrame(rowRDD, structType);
// 4.注册为临时表
personsDF.registerTempTable("test");
DataFrame result = hiveContext.sql("select * from test");
/**
* 对结果进行处理,包括由DataFrame转换成为RDD,以及结果的持久化
*/
List listRow = result.javaRDD().collect();
for (Row row : listRow) {
logger.error("row:" + row);
}
hiveContext.sql("insert into recommendation_system.t111 select id from test");
}
});
// 测试流,需要存在感~
lines.flatMap(new FlatMapFunction() {
public Iterable call(String msg) {
System.err.println(msg);
logger.error(msg);
return Lists.newArrayList(" ".split(msg));
}
}).mapToPair(new PairFunction() {
@Override
public Tuple2 call(String t) throws Exception {
return new Tuple2(t, 1);
}
}).reduceByKey(new Function2() {
@Override
public Integer call(Integer v1, Integer v2) throws Exception {
return v1 + v2;
}
}).print();
jssc.start();
jssc.awaitTermination();
}
}