Docker环境下Hive与Kafka连接问题及Hive3的docker-compose配置需求
基于Docker Compose的Hive 3 + Kafka环境配置方案
以下是适配Hive 3.1.2的完整docker-compose配置,同时修复了原配置中的端口映射错误、组件兼容性问题,确保Kafka与Hive生态的连通性:
version: "3" services: namenode: image: bde2020/hadoop-namenode:1.1.0-hadoop2.8-java8 container_name: namenode volumes: - namenode:/hadoop/dfs/name - ./infra/zeppelin/examples:/opt/sansa-examples environment: - CLUSTER_NAME=test env_file: - ./infra/hadoop/hadoop-hive.env ports: - "50070:50070" - "8020:8020" - "8081:8081" datanode: image: bde2020/hadoop-datanode:1.1.0-hadoop2.8-java8 container_name: datanode volumes: - datanode:/hadoop/dfs/data env_file: - ./infra/hadoop/hadoop-hive.env depends_on: - namenode spark-master: image: bde2020/spark-master:3.0.1-hadoop2.8-hive3.1.2 container_name: spark-master ports: - "8080:8080" - "7077:7077" environment: - CORE_CONF_fs_defaultFS=hdfs://namenode:8020 - SPARK_PUBLIC_DNS=localhost depends_on: - namenode - datanode - hive-metastore spark-worker: image: bde2020/spark-worker:3.0.1-hadoop2.8-hive3.1.2 container_name: spark-worker ports: - "8083:8083" environment: - SPARK_MASTER=spark://spark-master:7077 - CORE_CONF_fs_defaultFS=hdfs://namenode:8020 - SPARK_PUBLIC_DNS=localhost depends_on: - spark-master hue: image: bde2020/hdfs-filebrowser:3.11 container_name: hue ports: - 8088:8088 environment: - NAMENODE_HOST=namenode - SPARK_MASTER=spark://spark-master:7077 depends_on: - spark-master zeppelin: image: bde2020/zeppelin:0.9.0-hadoop2.8-spark3.0-hive3.1.2 container_name: zeppelin ports: - 8080:8080 volumes: - ./data:/data - ./data:/opt/zeppelin/data - ./infra/zeppelin/logs:/opt/zeppelin/logs - ./infra/zeppelin/notebooks:/opt/zeppelin/notebook - ./infra/zeppelin/examples:/opt/sansa-examples environment: CORE_CONF_fs_defaultFS: "hdfs://namenode:8020" SPARK_MASTER: "spark://spark-master:7077" MASTER: "spark://spark-master:7077" SPARK_SUBMIT_OPTIONS: "--jars /opt/sansa-examples/jars/sansa-examples-spark.jar --conf spark.serializer=org.apache.spark.serializer.KryoSerializer" depends_on: - spark-master # Hive 3.1.2 服务配置 hive-server: image: apache/hive:3.1.2 container_name: hive-server env_file: - ./infra/hadoop/hadoop-hive.env environment: - HIVE_SITE_CONF_javax_jdo_option_ConnectionURL=jdbc:postgresql://hive-metastore-postgresql:5432/metastore - HIVE_SITE_CONF_javax_jdo_option_ConnectionDriverName=org.postgresql.Driver - HIVE_SITE_CONF_javax_jdo_option_ConnectionUserName=hive - HIVE_SITE_CONF_javax_jdo_option_ConnectionPassword=hive - HIVE_SITE_CONF_datanucleus_autoCreateSchema=false - HIVE_SITE_CONF_hive_metastore_uris=thrift://hive-metastore:9083 - HADOOP_CONF_DIR=/opt/hadoop/etc/hadoop volumes: - ./infra/hadoop/hadoop-conf:/opt/hadoop/etc/hadoop ports: - 10000:10000 - 10002:10002 depends_on: - namenode - hive-metastore - hive-metastore-postgresql hive-metastore-postgresql: image: postgres:11 container_name: hive-metastore-postgresql environment: - POSTGRES_USER=hive - POSTGRES_PASSWORD=hive - POSTGRES_DB=metastore volumes: - hive-metastore-db:/var/lib/postgresql/data hive-metastore: image: apache/hive:3.1.2 container_name: hive-metastore env_file: - ./infra/hadoop/hadoop-hive.env environment: - HIVE_SITE_CONF_javax_jdo_option_ConnectionURL=jdbc:postgresql://hive-metastore-postgresql:5432/metastore - HIVE_SITE_CONF_javax_jdo_option_ConnectionDriverName=org.postgresql.Driver - HIVE_SITE_CONF_javax_jdo_option_ConnectionUserName=hive - HIVE_SITE_CONF_javax_jdo_option_ConnectionPassword=hive - HIVE_SITE_CONF_datanucleus_autoCreateSchema=true - HADOOP_CONF_DIR=/opt/hadoop/etc/hadoop volumes: - ./infra/hadoop/hadoop-conf:/opt/hadoop/etc/hadoop command: /opt/hive/bin/hive --service metastore ports: - 9083:9083 depends_on: - namenode - hive-metastore-postgresql zookeeper: image: confluentinc/cp-zookeeper:7.3.0 container_name: zookeeper environment: ZOOKEEPER_CLIENT_PORT: 2181 ZOOKEEPER_TICK_TIME: 2000 volumes: - zookeeper:/var/lib/zookeeper/data - zookeeper-logs:/var/lib/zookeeper/log kafka: image: confluentinc/cp-kafka:7.3.0 container_name: kafka ports: - "9092:9092" - "29092:29092" environment: KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181 KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9092,PLAINTEXT_HOST://localhost:29092 KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1 volumes: - kafka:/var/lib/kafka/data depends_on: - zookeeper nifi: image: apache/nifi:1.19.0 container_name: nifi ports: - 5080:8080 - 5443:8443 environment: - NIFI_WEB_HTTP_PORT=8080 - NIFI_CLUSTER_IS_NODE=true - NIFI_CLUSTER_NODE_PROTOCOL_PORT=8082 - NIFI_ZK_CONNECT_STRING=zookeeper:2181 - NIFI_ELECTION_MAX_WAIT=1 min volumes: - ./infra/nifi/conf:/opt/nifi/nifi-current/conf - ./infra/nifi/logs:/opt/nifi/nifi-current/logs - ./data:/opt/datafiles - nifi-flowfile:/opt/nifi/nifi-current/flowfile_repository - nifi-db:/opt/nifi/nifi-current/database_repository - nifi-content:/opt/nifi/nifi-current/content_repository - nifi-provenance:/opt/nifi/nifi-current/provenance_repository depends_on: - zookeeper - kafka volumes: namenode: datanode: zookeeper: zookeeper-logs: kafka: nifi-flowfile: nifi-db: nifi-content: nifi-provenance: hive-metastore-db:
关键改动说明
- Hive版本升级:替换为
apache/hive:3.1.2官方镜像,搭配PostgreSQL 11作为元数据存储,确保Hive 3的特性支持 - Spark版本适配:使用支持Hive 3的Spark 3.0.1镜像,保证Spark与Hive的兼容性
- Kafka配置优化:改用Confluent Kafka镜像,配置双监听地址(容器内部+宿主机),避免网络连通问题
- 端口映射修复:修正原spark-master的错误端口映射(原8090:800改为8080:8080)
- 依赖关系调整:明确各服务的依赖顺序,确保元数据存储、HDFS等核心服务先启动
Kafka数据写入Hive的实现方式
方式1:Spark Structured Streaming
在Zeppelin中执行以下Scala代码,消费Kafka Topic并写入Hive表:
import org.apache.spark.sql.streaming.Trigger // 初始化Hive支持 spark.sql("CREATE DATABASE IF NOT EXISTS kafka_db") spark.sql("USE kafka_db") spark.sql(""" CREATE TABLE IF NOT EXISTS kafka_metadata ( key STRING, value STRING, topic STRING, partition INT, offset BIGINT, timestamp TIMESTAMP ) STORED AS PARQUET """) // 消费Kafka数据 val kafkaDF = spark.readStream .format("kafka") .option("kafka.bootstrap.servers", "kafka:9092") .option("subscribe", "your_topic_name") .load() // 转换数据结构 val parsedDF = kafkaDF.selectExpr( "CAST(key AS STRING)", "CAST(value AS STRING)", "topic", "partition", "offset", "timestamp" ) // 写入Hive表 val query = parsedDF.writeStream .format("hive") .trigger(Trigger.ProcessingTime("10 seconds")) .option("checkpointLocation", "/tmp/checkpoint/kafka_metadata") .table("kafka_metadata") query.awaitTermination()
方式2:Apache NiFi
- 添加
ConsumeKafkaRecord_2_6处理器,配置Kafka bootstrap地址为kafka:9092,指定要消费的Topic - 添加
ConvertRecord处理器,将Kafka的二进制数据转换为JSON/CSV格式 - 添加
PutHiveStreaming处理器,配置Hive Metastore地址为thrift://hive-metastore:9083,指定目标Hive表 - 连接各处理器,启动数据流
内容的提问来源于stack exchange,提问作者ib 1937
相关产品推荐
相关产品推荐

