Showing posts with label Java. Show all posts
Showing posts with label Java. Show all posts
Monday, November 9, 2015
Monday, October 19, 2015
Thursday, September 17, 2015
ReactiveX in Java
Dependency:
<dependency><groupId>io.reactivex</groupId>
<artifactId>rxjava</artifactId>
<version>1.0.14</version>
</dependency>
Code:
Observable:ThreadPoolExecutor executor = new ThreadPoolExecutor(3, 5, 3000L, TimeUnit.MILLISECONDS, new LinkedBlockingQueue());
Observable<ChannelItemType> myObservable = Observable.create(
new Observable.OnSubscribe<ChannelItemType>() {
@Override
public void call(Subscriber<? super ChannelItemType> sub) {
sub.onNext(item);
sub.onCompleted();
}
}
).subscribeOn(Schedulers.from(executor));
myObservable.subscribe(new ChannelConsumer(channelRepo));
Above snippet, we create an Observable, sent item on onNext() function to Consumer (in this case ChannelConsumer). One interesting thing here is we thread consuming by using executor. We create 3 thread in executor. Obervable automatically pick thread in executor and send data.
Consumer:
public class ChannelConsumer extends Subscriber<ChannelItemType> {
@Override
public void onCompleted() {
//Clean up your process when complete
}
@Override
public void onError(Throwable throwable) {
}
@Override
public void onNext(ChannelItemType channelItem) {
//put your code here
}
}
Monday, August 17, 2015
Batch insert with JDBC and Spring
JDBC:
String sql = "insert into employee (name, city, phone) values (?, ?, ?)";
Connection connection = new getConnection();
PreparedStatement ps = connection.prepareStatement(sql);
final int batchSize = 1000;
int count = 0;
for (Employee employee: employees) {
ps.setString(1, employee.getName());
ps.setString(2, employee.getCity());
ps.setString(3, employee.getPhone());
ps.addBatch();
if(++count % batchSize == 0) {
ps.executeBatch();
}
}
ps.executeBatch(); // insert remaining records
ps.close();
connection.close();
Connection connection = new getConnection();
PreparedStatement ps = connection.prepareStatement(sql);
final int batchSize = 1000;
int count = 0;
for (Employee employee: employees) {
ps.setString(1, employee.getName());
ps.setString(2, employee.getCity());
ps.setString(3, employee.getPhone());
ps.addBatch();
if(++count % batchSize == 0) {
ps.executeBatch();
}
}
ps.executeBatch(); // insert remaining records
ps.close();
connection.close();
Spring Data:
public class EmployeeRepository extends JdbcDaoSupport {
public void insertEmployeeBatch(List<Employee> employees) {
try {
String sql = "insert into employee (name, city, phone) values (?, ?, ?)";
getJdbcTemplate().batchUpdate(sql, new BatchPreparedStatementSetter() {
public void setValues(PreparedStatement ps, int i) throws SQLException { Employee employee = employees.get(i); ps.setString(1, employee.getName());
ps.setString(2, employee.getCity());
ps.setString(3, employee.getPhone());
}
public int getBatchSize() {
return employees.size();
}
});
} catch (Exception ex) {
LOGGER.error(ex.getMessage(), ex);
}
} }Thursday, August 13, 2015
Read big XML file in Java
Tool:
LTFViewerMaven:
<dependency><groupId>org.codehaus.woodstox</groupId>
<artifactId>stax2-api</artifactId>
<version>3.1.2</version>
</dependency>
Java:
public void parseFile(String filename, Long jobId) throws Exception {XMLStreamReader2 xmlr = null;
FileInputStream is = null;
List<Future> trackingTask = null;
try {
XMLInputFactory2 xmlif = ((XMLInputFactory2) XMLInputFactory.newInstance());
xmlif.configureForSpeed();
// Start init executor service
trackingTask = initExecutorService();
is = new FileInputStream(filename);
xmlr = (XMLStreamReader2) xmlif.createXMLStreamReader(is);
// Parse into typed objects
JAXBContext ctx = JAXBContext.newInstance(PGWSSchedulePOJO.class, PGWSChannelPOJO.class...);
Unmarshaller um = ctx.createUnmarshaller();
while (xmlr.hasNext()) {
xmlr.next();
if (xmlr.isStartElement()) {
if ((xmlr.getLocalName().equals("publishedTitles"))) {
NlpgwsTitleInfoPOJO ti = um.unmarshal(xmlr, NlpgwsTitleInfoPOJO.class).getValue();
if (ti != null) {
//Add to queue "parsedTitleInfo"
parsedTitleInfo.add(ti);
}
} else if ((xmlr.getLocalName().equals("providers"))) {
NlpgwsProviderPOJO pr = um.unmarshal(xmlr, NlpgwsProviderPOJO.class).getValue();
if (pr != null) {
//Add to queue "parsedProviders"
parsedProviders.add(pr);
}
}
}
}
} catch (XMLStreamException ex) {
LOGGER.error(ex.getMessage(), ex);
throw ex;
} catch (Exception ex) {
LOGGER.error(ex.getMessage(), ex);
throw ex;
} finally {
if (xmlr != null) {
try {
xmlr.close();
} catch (Exception ex) {
LOGGER.error(ex.getMessage());
}
}
if (is != null) {
try {
is.close();
} catch (Exception ex) {
LOGGER.error(ex.getMessage());
}
}
}
// Waiting for all task to completed and shutdown the executor server
try {
if (trackingTask != null) {
for (Future future : trackingTask) {
future.get();
}
}
} catch (Exception ex) {
LOGGER.error(ex.getMessage(), ex);
} finally {
// clean job
}
}
With:
protected static BlockingQueue<NlpgwsProviderPOJO> parsedProviders = new LinkedBlockingDeque<NlpgwsProviderPOJO>();
protected static BlockingQueue<NlpgwsTitleInfoPOJO> parsedTitleInfo = new LinkedBlockingDeque<NlpgwsTitleInfoPOJO>();
Algorithm here:
- Read XML file
- Put data object to a BlockingQueue
- Init Executor -> create thread to scan the Blocking queue to do the business logic
Declare executor:
protected ExecutorService ingestionBatchInsert;
ThreadFactory batchInsertThreadFactory = new ThreadFactoryBuilder().setNameFormat("dls-batch-insert-%d").build();
ingestionBatchInsert = Executors.newFixedThreadPool(totalInsertThread + 1, batchInsertThreadFactory);
Add thread in trackingList:
for (int count = 0; count < totalInsertThread; count++) {
Future task = ingestionBatchInsert.submit(othersBatchInsertThread);
trackingTask.add(task);
}
OthersBatchInsertThread thread:
public class OthersBatchInsertThread implements Runnable {
public void run() {
//Get out object from BlockingQueue
DLSFileParserService.parsedProviders.poll()
//process
}
}
Check file is open (lock) or not in Java
RandomAccessFile rf = new RandomAccessFile(file, "rw");
FileChannel fileChannel = rf.getChannel();
FileLock lock = null;
try {
// let us try to get a lock. If file already has an exclusive lock by another process
LOGGER.info("Trying to acquire lock");
lock = fileChannel.tryLock();
if (lock != null) {
success = true;
}
} catch (Exception ex) {
LOGGER.error(ex.getMessage());
} finally {
if (lock != null) {
lock.release();
}
if(fileChannel != null){
fileChannel.close();
}
if(rf != null){
rf.close();
}
}
FileChannel fileChannel = rf.getChannel();
FileLock lock = null;
try {
// let us try to get a lock. If file already has an exclusive lock by another process
LOGGER.info("Trying to acquire lock");
lock = fileChannel.tryLock();
if (lock != null) {
success = true;
}
} catch (Exception ex) {
LOGGER.error(ex.getMessage());
} finally {
if (lock != null) {
lock.release();
}
if(fileChannel != null){
fileChannel.close();
}
if(rf != null){
rf.close();
}
}
Observe folder to pick file when it's available
public class MMTServerStartListener implements ApplicationListener<ContextRefreshedEvent> {
@Autowired
YourMonitorListenerImpl fileMonitor;
@Override
public void onApplicationEvent(ContextRefreshedEvent context) {
try {
String filePath = "<file path>";
startMonitor(filePath, fileMonitor);
} catch (Exception ex) {
LOGGER.error(ex.getMessage(), ex);
}
}
private void startMonitor(String filePath, YourMonitorListenerImpl fileMonitorImpl) {
try {
if (filePath != null && filePath.length() > 0) {
final File directory = new File(filePath.trim());
FileAlterationObserver fao = new FileAlterationObserver(directory);
fao.addListener(fileMonitorImpl);
final FileAlterationMonitor monitor = new FileAlterationMonitor();
monitor.addObserver(fao);
LOGGER.info("Starting monitor. CTRL+C to stop.");
monitor.start();
Runtime.getRuntime().addShutdownHook(new Thread(new Runnable() {
public void run() {
try {
LOGGER.info("Stopping monitor.");
monitor.stop();
} catch (Exception ignored) {
LOGGER.error(ignored.getMessage(), ignored);
}
}
}));
} else {
LOGGER.error("Invalid input the monitor folder");
}
} catch (Exception ex) {
LOGGER.error(ex.getMessage(), ex);
}
}
}
With "YourMonitorListenerImpl" implements "FileAlterationListener" interface.
Wednesday, August 12, 2015
Thread with ExecutorService
ThreadFactory buildCache = new ThreadFactoryBuilder().setNameFormat("thread-name-%d").build();
ExecutorService executorService = Executors.newFixedThreadPool(totalThreads, buildCache);
List<Future> trackingTask = new ArrayList<Future>();
for (int index = 0; index < totalThreads; index++) {
trackingTask.add(executorService.submit(new Runnable() {
@Override
public void run() {
//run your code
}
}));
}
// Run and waiting for task to finish
for (Future task : trackingTask) {
task.get();
}
// Shutdown executor:
executorService.shutdownNow();
ExecutorService executorService = Executors.newFixedThreadPool(totalThreads, buildCache);
List<Future> trackingTask = new ArrayList<Future>();
for (int index = 0; index < totalThreads; index++) {
trackingTask.add(executorService.submit(new Runnable() {
@Override
public void run() {
//run your code
}
}));
}
// Run and waiting for task to finish
for (Future task : trackingTask) {
task.get();
}
// Shutdown executor:
executorService.shutdownNow();
Tuesday, August 11, 2015
Synchronize SimpleDateFormat object in Java
SimpleDateFormat object does not work properly in a multi threaded environment. It may output a wrong date when parsing. So the safest way is to synchronize it.
private SimpleDateFormat sdf = new SimpleDateFormat("yyyy-MM-dd");
public String formatDate(Date d) {
synchronized(sdf) {
return sdf.format(d);
}
}
Hope that help.
Sunday, August 9, 2015
Check if a class is loaded and lib file location in JVM
Example below is to see if JVM using ojdbc6 version 11 or version 12.
final ClassLoader loader = Thread.currentThread().getContextClassLoader();
ClassPath clazzPath = ClassPath.from(loader);
Set<ClassInfo> classes = clazzPath.getTopLevelClasses();
for (final ClassPath.ClassInfo classInfo : classes) {
if (classInfo.getName().startsWith("oracle")) {
if (classInfo.getName().contains("oracle.jdbc.babelfish")) {
isVersion12 = true;
}
}
}
if(isVersion12) {
LOGGER.info("Ojdbc version 12");
} else {
LOGGER.info("Ojdbc version 11");
}
Class klass = OracleConnection.class;
URL location = klass.getResource('/' + klass.getName().replace('.', '/') + ".class");
LOGGER.info(location.toString());
final ClassLoader loader = Thread.currentThread().getContextClassLoader();
ClassPath clazzPath = ClassPath.from(loader);
Set<ClassInfo> classes = clazzPath.getTopLevelClasses();
for (final ClassPath.ClassInfo classInfo : classes) {
if (classInfo.getName().startsWith("oracle")) {
if (classInfo.getName().contains("oracle.jdbc.babelfish")) {
isVersion12 = true;
}
}
}
if(isVersion12) {
LOGGER.info("Ojdbc version 12");
} else {
LOGGER.info("Ojdbc version 11");
}
Class klass = OracleConnection.class;
URL location = klass.getResource('/' + klass.getName().replace('.', '/') + ".class");
LOGGER.info(location.toString());
Read package name and version from Manifest file
Bean:
public class BuildVersion {/** The build name. */
private String buildName;
/** The version. */
private String version;
/**
* Gets the version.
*
* @return the version
*/
public String getVersion() {
return version;
}
/**
* set version.
*
* @param version the new version
*/
public void setVersion(String version) {
this.version = version;
}
/**
* Gets the builds the name.
*
* @return the builds the name
*/
public String getBuildName() {
return buildName;
}
/**
* Sets the builds the name.
*
* @param buildName the new builds the name
*/
public void setBuildName(String buildName) {
this.buildName = buildName;
}
}
VersionServiceImpl:
private BuildVersion getBuildVersion() {BuildVersion buildVersion = new BuildVersion();
Class clazz = this.getClass();
String className = clazz.getSimpleName() + ".class";
String classPath = clazz.getResource(className).toString();
try {
String manifestPath = classPath.substring(0, classPath.lastIndexOf("/WEB-INF")) +
"/META-INF/MANIFEST.MF";
Manifest manifest = new Manifest(new URL(manifestPath).openStream());
Attributes attr = manifest.getMainAttributes();
String title = attr.getValue("Specification-Title");
String version = attr.getValue("Specification-Version");
buildVersion.setBuildName(title);
buildVersion.setVersion(version);
} catch (Exception ex) {
LOGGER.error(ex.getMessage(), ex);
}
return buildVersion;
}
VersionEndpoint:
@Component@Path ("/")
public class VersionEndpoint {
/** The version service. */
@Autowired
private VersionService versionService;
/**
* getVersion service use to retrieve the current release version.
*
* @param format the format
* @return Version response
*/
@GET
@Path ("/version")
@Produces ({ MediaType.APPLICATION_XML, MediaType.APPLICATION_JSON })
public Response getVersion(@DefaultValue (WSConstants.RESPONSE_TYPE_JSON) @QueryParam ("format") String format) {
try {
return Response.ok().entity(versionService.getVersion()).type(WSUtil.getWSResponseType(format)).build();
} catch (Exception e) {
return Response.ok().entity("Error: " + e.getMessage()).type(WSUtil.getWSResponseType(format)).build();
}
}
}
Tuesday, January 13, 2015
Deeper understand of Enum in Java
Reference:http://howtodoinjava.com/2012/12/07/guide-for-understanding-enum-in-java/ Enumerations (in general) are generally a set of related constants. Each constant represents a specific instance of its Enumerations class. For example:
public enum DIRECTION { EAST, WEST, NORTH, SOUTH //optionally can end with ";" } Here EAST, WEST, NORTH and SOUTH are final static inner classes of Direction of type Direction extends java.lang.Enum.Enums are comparable and serializable implicitly. Also, all enum types in java are singleton by default. So, you can compare enum types using ‘==’ operator.
You can give define your own constructors to initialize the state of enum types.enum Direction { // Enum types EAST(0), WEST(180), NORTH(90), SOUTH(270); // Constructor private Direction(final int angle) { this.angle = angle; } // Internal state private int angle; public int getAngle() { return angle; }}
Use template methods in enum:
package enumTest;public enum Direction { // Enum types EAST(0) { <strong>@Override</strong> public void shout() { System.out.println("Direction is East !!!"); } }, WEST(180) { @Override public void shout() { System.out.println("Direction is West !!!"); } }, NORTH(90) { @Override public void shout() { System.out.println("Direction is North !!!"); } }, SOUTH(270) { @Override public void shout() { System.out.println("Direction is South !!!"); } }; // Constructor private Direction(final int angle) { this.angle = angle; } // Internal state private int angle; public int getAngle() { return angle; } <strong>// Abstract method which need to be implemented</strong> public abstract void shout();}
Subscribe to:
Posts (Atom)