在Flink SQL中读取Kafka数据并写入Hive表时,要确保数据的类型和字段顺序与Hive表的定义一致。根据你提供的样例数据和Hive表定义,可以通过以下步骤来实现:

  1. 创建Flink SQL的执行环境。
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
StreamTableEnvironment tEnv = StreamTableEnvironment.create(env);
  1. 注册Kafka表。
String kafkaBootstrapServers = "localhost:9092";
String kafkaTopic = "your_topic";
String groupId = "your_group_id";

String createKafkaSource = String.format("CREATE TABLE kafka_source (\n" +
        "  roadGBCode STRING,\n" +
        "  totalFlow DOUBLE,\n" +
        "  storageTime BIGINT,\n" +
        "  occupancy DOUBLE,\n" +
        "  roadId STRING,\n" +
        "  startMilestone INT,\n" +
        "  fictitious_road_code STRING,\n" +
        "  uuid STRING,\n" +
        "  speed85 DOUBLE,\n" +
        "  endMilestone INT,\n" +
        "  congestionLevel INT,\n" +
        "  gap INT,\n" +
        "  start_hour STRING,\n" +
        "  id INT,\n" +
        "  lane INT,\n" +
        "  direction INT,\n" +
        "  timestamp BIGINT,\n" +
        "  smallCarQuantity INT,\n" +
        "  headway INT,\n" +
        "  averageSpeed DOUBLE,\n" +
        "  totalSpeed DOUBLE,\n" +
        "  largeCarQuantity INT,\n" +
        "  fictitious_milestone STRING,\n" +
        "  mediumCarQuantity INT,\n" +
        "  circle INT,\n" +
        "  vehicleQuantity INT\n" +
        ") WITH (\n" +
        "  'connector' = 'kafka',\n" +
        "  'topic' = '%s',\n" +
        "  'properties.bootstrap.servers' = '%s',\n" +
        "  'properties.group.id' = '%s',\n" +
        "  'scan.startup.mode' = 'earliest-offset',\n" +
        "  'format' = 'json',\n" +
        "  'json.fail-on-missing-field' = 'false'\n" +
        ")", kafkaTopic, kafkaBootstrapServers, groupId);

tEnv.executeSql(createKafkaSource);
  1. 创建Hive表。
String createHiveTable = "CREATE TABLE hive_table (\n" +
        "  start_milestone INT COMMENT '起始断面位置桩号',\n" +
        "  end_milestone INT COMMENT '路段结束桩号',\n" +
        "  timestamp BIGINT COMMENT '起始时间',\n" +
        "  lane INT COMMENT '车道号',\n" +
        "  direction INT COMMENT '路段方向',\n" +
        "  vehicle_quantity INT COMMENT '车道流量',\n" +
        "  average_speed DOUBLE COMMENT '车道速度',\n" +
        "  total_flow DOUBLE COMMENT '路段流量'\n" +
        ")\n" +
        "COMMENT 'kafka2hive 实时测试'\n" +
        "PARTITIONED BY (ds STRING COMMENT '分区')\n" +
        "STORED AS PARQUET\n" +
        "LOCATION 'hdfs://your_hdfs_path'";

tEnv.executeSql(createHiveTable);

注意:你需要将your_hdfs_path替换为实际的HDFS路径。

  1. 将Kafka数据写入Hive表。
String insertIntoHive = "INSERT INTO hive_table\n" +
        "SELECT\n" +
        "  startMilestone AS start_milestone,\n" +
        "  endMilestone AS end_milestone,\n" +
        "  timestamp,\n" +
        "  lane,\n" +
        "  direction,\n" +
        "  vehicleQuantity AS vehicle_quantity,\n" +
        "  averageSpeed AS average_speed,\n" +
        "  totalFlow AS total_flow\n" +
        "FROM kafka_source";

tEnv.executeSql(insertIntoHive);
  1. 启动Flink任务。
env.execute();

以上代码将会从Kafka读取数据,并将指定字段映射到Hive表中。注意要将your_topicyour_group_id和HDFS路径替换为实际的值。另外,如果你想要指定特定的分区值,可以在插入语句中使用PARTITION子句来指定分区的值。

flinksql 读取kafka数据写入hivekafka的样例数据为roadGBCodeG60totalFlow280storageTime1694677876000occupancy-10roadId86400startMilestone86400fictitious_road_code1001uuid799760-1186500-1-1311-10G608640046797864001694

原文地址: https://www.cveoy.top/t/topic/i8wB 著作权归作者所有。请勿转载和采集!

免费AI点我,无需注册和登录