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();

	}
}