logo
publist
写文章

简介

该用户还未填写简介

擅长的技术栈

可提供的服务

暂无可提供的服务

flink维表关联系列之kafka维表关联:广播方式

Flink中广播状态假设存在这样一种场景,一个是用户行为数据,一个是规则数据,要求通过规则去匹配用户行为找到符合规则的用户,并且规则是可以实时变更的,在用户行为匹配中也能根据规则的实时变更作出相应的调整。这个时候就可以使用广播状态,将用户行为数据看做是一个流userActionStream,规则数据也看做是一个流ruleStream,将ruleStream流中数据下发到userActionStre

#flink#kafka#大数据
flink 的consumer.setStartFromTimestamp方法

问题描述:在消费Kafka的时候通过consumer.setStartFromTimestamp方法将消费起始点设置到12月3号0点0分0秒,但是消费到的数据有12月2号23点59分59秒的。问题原因:consumer.setStartFromTimestamp 表示数据写进 kafka 的时间...

#flink#大数据#kafka
kafka connect结合debezium采集oracle数据的时候,任务失败重启的数据偏移量问题

这将导致旧 SCN 值和新提供的 SCN 值之间发生的更改丢失,并且不会写入主题。当连接器报告找不到此偏移 SCN 时,这表明仍然可用的日志不包含 SCN,因此连接器无法从它停止的地方挖掘更改。Debezium Oracle 连接器在偏移量中维护两个关键值,一个名为 scn 的字段 和另一个名为commit_scn的字段。找出连接器的最后一个偏移量、存储它的键并确定用于存储偏移量的分区。重启对应的

#kafka#oracle#java
FlinkSQL CDC实现同步oracle数据到mysql

环境准备1、flink 1.13.02、oracle 11g3、flink-connector-oracle-cdc 2.1.01、oracle环境配置首先需要安装oracle环境,参考 https://blog.csdn.net/qq_36039236/article/details/124224500?spm=1001.2014.3001.5502进入容器进行配置:docker exec -i

#flink#oracle
hive sql - 如何选择 hive 数组列中的前 n 个元素并返回选定的数组

表数据如下:user_idinterest_arraytom[a,b,c,d,g,w]bob[e,d,s,d,g,w,s]cat[a]harry[]peterNULL目的是按顺序选择每行“interest_array”中的前 3 个元素并将其作为数组返回,输出如下:user_idoutput_arraytom[a,b,c]bob[e,d,s]cat[a]harry[]peter

#hive#sql
impala: 类似hive的array_contains()函数

impala: 类似hive的array_contains()函数首先,需要先了解impala中的两个函数的作用,一个是group_concat(string s [, string sep]), 一个是find_in_set(string str,string strList)1、group_concat(string s [, string sep])按照指定分隔符, 将多行记录的 s 表达式

#大数据
flink sql在实时数仓中,关联hbase维表频繁变化的问题

在用flink sql在做实时数仓,架构大概是kafka关联hbase维表,然后写入clickhouse。hbase维表是频繁变化的现在遇到的几个比较棘手的问题:1、自己在实现AsyncTableFunction做异步io的时候,发现性能还是不够。后来就加入本地缓存,但是缓存一致性出现问题,不知道该如何解决2、写入hbase的时候,是批量写的,无法保证有序,维表频繁变化的话,顺序不对,会造成结果有

#flink#大数据
hive : 使用str_to_map 和regexp_replace , 将 json 字符串转换为 map

str_to_map(字符串参数, 分隔符1, 分隔符2)使用两个分隔符将文本拆分为键值对。分隔符1将文本分成K-V对,分隔符2分割每个K-V对。对于分隔符1默认分隔符是 ‘,’,对于分隔符2默认分隔符是 ‘=’。使用案例:select str_to_map('aaa:11&bbb:22', '&', ':')结果:{"bbb":"22","aaa":"11"}将 json 字符串

#hive#大数据
flink维表关联系列之Mysql维表关联:全量加载

在维表关联中定时全量加载是针对维表数据量较少并且业务对维表数据变化的敏感程度较低的情况下可采取的一种策略,对于这种方案使用有几点需要注意:全量加载有可能会比较耗时,所以必须是一个异步加载过程内存维表数据需要被流表数据关联读取、也需要被定时重新加载,这两个过程是不同线程执行,为了尽可能保证数据一致性,可使用原子引用变量包装内存维表数据对象,即AtomicReference查内存维表数据非异步io过程

#flink#大数据
hive 中笛卡尔积的优化 (大表/小表)

目录笛卡尔积处理测试案例笛卡尔积处理当Hive设定为严格模式(hive.mapred.mode = strict)时,不允许在HQL语句中出现笛卡尔积,这实际说明了Hive 对笛卡尔积支持较弱。因为找不到 join key, Hive只能使用一个reducer 来完成笛卡尔积。当然也可以使用 limit 的办法来减少某个表参与 join 的数据量,但对于需要笛卡尔积语义的需求来说,经常是一个大表和

#大数据#hive
    共 20 条
  • 1
  • 2
  • 请选择