flink相关 – 读入kafka数据源

业务代码如下 :
fsTableEnv.connect(
       newKafka()
                .version("0.11")
                .topic("xxxx")
                .property("bootstrap.servers”,”xxxxx:xxxx")
                .property("group.id","test4444")
                .startFromEarliest()//测试需要
)
        .withSchema(newSchema().schema(newTableSchema(fields, types)))
        .withFormat(newJson().schema(newRowTypeInfo(types, fields)))
        .inAppendMode()
        .registerTableSource(“testTable”);
问题是:读出来的数据总是null
DataStream table = fsTableEnv.toAppendStream(fsTableEnv.sqlQuery("select * from testTable"),Row.class);
table.print();

继续阅读“flink相关 – 读入kafka数据源”

flink学习笔记(一)

架构
  1.  什么是flink
flink是一个框架或者叫流式计算引擎,他可以用来处理有界(批处理,mapReduce)数据和无界数据。他的特点是具有内存的速度和分布式的横向扩展能力。由于是分布式的,所以他可以处理任意规模的数据。另外,他是有状态的,表现为上一步的结算结果可以为下一步使用。

继续阅读“flink学习笔记(一)”

JVM Sandbox 源码分析(一)基础篇: Instrumentation和ClassFileTransformer作用

背景

jvm sandbox是一个开源项目https://github.com/alibaba/jvm-sandbox/wiki/FIRST-MODULE,他是对java agent的一个应用,他将agent发扬光大,让你可以不用用户改代码的情况下,进行代码增强,实现AOP,切面编程的功能,非常强大。但是他的代码量其实不大,文档清晰,所以读起来不是太费劲。他最最基础和核心的功能是使用了Java 提供的 Instrumenttation和ClassFileTransformer来实现的,我们先来了解下这个两个类的API,了解他们都提供了什么功能。

 

继续阅读“JVM Sandbox 源码分析(一)基础篇: Instrumentation和ClassFileTransformer作用”

Java类库设计之日志框架选择

Java日志框架现状

Java日志框架的故事说来话长,做过开发的一定遇到过slf4j,log4j,logback,commons-logging,log4j2。=也一定见过这些jar包,什么log4j.jar,slf4j-api.jar,slf4j-log4j12.jar,logback-classic.jar,logback-core.jar等等。以前看到他们就头大,傻傻分不清楚,到现在搞清楚了也就这么回事,可以具体看看slf4j的实现和适配是怎么做的。

继续阅读“Java类库设计之日志框架选择”

Java多线程知识点(下)

26) 如何写代码来解决生产者消费者问题?

我们先尝试使用队列实现一个生产者和消费者模式。那么就需要先了解队列,队里的接口定义如下
    Queue.
         添加:
                  add(E element) 添加一个,如果超出了界限,则抛出异常
                  Offer(E element) 添加一个,如果超出了界限,则返回false
        获取并删除
                  remove() 获取并删除队头,如果队列为空,则抛出异常
                  Poll() 获取并删除队头, 如果队列为空,则返回null
        获取但不删除
                   element() 获取但是并不删除,如果队列为空,则抛出异常
                  peek() 获取但是并不删除,如果队列为空,则返回null.
 而阻塞队列则多了两个阻塞的方法,put 和 take,参见下表

继续阅读“Java多线程知识点(下)”

如何写代码来解决生产者消费者问题?

这个题目考验的实际上是线程间的通信和同步问题。我们再有一篇文章中,使用Object的wait和notify实现了线程间的同步。参见 : https://huster.top/htmls/615.html

继续阅读“如何写代码来解决生产者消费者问题?”

Java多线程总结(上部分)面试题解答

    自去年列下计划到现在一年有余,看了好几本书,以为已经过了很长时间,现在回头发现仅仅才过了一年而已,已经积累了不少知识,看来有目的的去学习还是有一些成效的。正好今天看到一个关于多线程的面试题,所以来总结过去所学,也是对过去知识的一个考验。

继续阅读“Java多线程总结(上部分)面试题解答”

细说Java的内存模型

    多线程之所以复杂,是因为在多线程环境下,稍有不慎就会导致诡异的问题。这个诡异表现在,他不是一定会出现,他不知道会在什么时候会以什么样的姿态出现在我们的生产环境中。为了避免犯错误,所以我们在所有可能出问题的地方都加上synchronize,那当然也是欠妥的,所以我们要搞清楚其中的细节,对我们写程序就有比较大的帮助了。

继续阅读“细说Java的内存模型”