4 Haziran 2026 Perşembe

HazelcastAPI CP SubSystem vs Transaciton Yapılar

Giriş

Bu yazının yazılma sebebi bu soru

1. Transactional Yapılar

Hazelcast TransactionalMap, birden fazla veri değişikliğinin tek bir işlem (transaction) olarak atomik şekilde yapılmasını sağlar.
A = 100 -> 50
B = 50  -> 100
Transfer sırasında hata olursa tüm değişiklikler geri alınır.
Ancak bu, cluster'daki tüm düğümlerin aynı anda aynı veriyi gördüğü anlamına gelmez. Ağ bölünmeleri (split-brain), replikasyon gecikmeleri veya node arızalarında bazı düğümler geçici olarak eski veriyi görebilir.

Amaç:
- Atomic commit
- Rollback
- Isolation

2. CP SubSystem

CP Subsystem'in amacı transaction değil, dağıtık tutarlılıktır.

Burada Raft konsensüs algoritması kullanılır ve tüm CP üyeleri işlemlerin sırasını aynı şekilde kabul eder.
Örneğin:
counter = 5
counter = 6
counter = 7
Bir istemci 7'yi gördükten sonra başka bir istemcinin tekrar 6 görmesi mümkün değildir.

Ağ bölünmesi olursa sistem gerekirse isteği reddeder veya bekletir; yanlış/stale veri döndürmez.

3. CP Map vs Transaction doğru ayrım

3.1 CP Map:

Tekil operasyonlar için global doğruluk ve sıra garantisi

- “Herkes aynı sonucu aynı sırada görür”
- rollback yok
- her operation consensus’tan geçer

3.2 TransactionalMap:

Transaction başlatıldığında Hazelcast bir transaction context oluşturur. Bu context cluster genelinde coordinator + participant nodes üzerinden yönetilir. Ama, Hangi node’ların dahil olacağı senin seçtiğin bir liste değildir.

Kim belirler?

Hazelcast bunu otomatik belirler: Veri hangi partition’larda ise O partition’ların owner / backup node’ları transaction’a dahil olur. Transaction coordinator bunu yönetir

Yani birden fazla operation’ı atomik grup yapar
- “Ya hepsi olur ya hiçbiri”
- ama global ordering / linearizability garantisi yok

Linearizability olmadığı için 
Transaction commit:
x = 6
Durum:
Node A → 6 görüyor
Node B → 5 (eski değer) görüyor
Bu mümkündür.

HazelcastAPI CP Subsystem CPMap Arayüzü

Örnek ver

28 Temmuz 2024 Pazar

THIRD-PARTY.txt Dosyası

Kullanılan harici kütüphanelerin sürümleri bu dosyada
Dosyanın yolu şöyle

hazelcast/licenses/THIRD-PARTY.txt

9 Temmuz 2024 Salı

3 Haziran 2024 Pazartesi

Partition Table

Giriş
Açıklaması şöyle
When you start your first member, a partition table is created within it. As you start additional members, that first member becomes the oldest member, also known as the master member and updates the partition table accordingly. This member periodically sends the partition table to all other members. This way, each member in the cluster is informed about any changes to partition ownership. The ownerships may be changed when, for example, a new member joins the cluster, or when a member leaves the cluster.

NOTE : If the master member goes down, the next oldest member sends the partition table information to the other ones.

22 Ocak 2024 Pazartesi

HazelcastInstanceAware Arayüzü

Giriş
Şu satırı dahil ederiz
import com.hazelcast.core.ClientSchemaService;
Bu sınıfı ilk defa burada gördüm. Açıklaması şöyle
Entry Processor should implement HazelcastInstanceAware interface. It will provide setter to instance. There is no necessity to direct injection of instance. Hz will do it by itself on its side once EP would be deserialised there.


3 Ocak 2024 Çarşamba

Hazelcast Jet JobTerminateRequestedException Sınıfı

Giriş
Şu satırı dahil ederiz 
import com.hazelcast.jet.impl.exception.JobTerminateRequestedException;
Jet işleri nedense exception ile sonlandırılıyor. Sonlandırma sebebi TerminationMode ile belirtiliyor

JobTerminateRequestedException iş Jet Engine tarafından iptal edilince fırlatılır
CancellationByUserException iş kullanıcı tarafından iptal edilince fırlatılır.

TerminationMode şöyle olabilir
  CANCEL_FORCEFUL : İş bitmiştir ve kapatılması gerekir
  CANCEL_GRACEFUL

  RESTART_FORCEFUL : İşin tekrar başlatılması gerekir
  RESTART_GRACEFUL : İşin tekrar başlatılması gerekir

 SUSPEND_FORCEFUL : İşin askıya alınması gerekir
 SUSPEND_GRACEFUL
  




8 Aralık 2023 Cuma

Socket Ayarları

Örnek
Şöyle yaparız
-Dhazelcast.socket.receive.buffer.size=10240 -Dhazelcast.socket.send.buffer.size=10240

7 Aralık 2023 Perşembe

Hazelcast SQL Sütun Tipleri

Sütun Tipleri burada

com.hazelcast.jet.sql.impl.connector.keyvalue.JavaClassNameResolver sütun tiplerinin hangi Java tipine karşılık geleceğini bilir.

6 Aralık 2023 Çarşamba

HazelcastAPI ClientSchemaService Sınıfı

Giriş
Şu satırı dahil ederiz
import com.hazelcast.client.impl.clientside.ClientSchemaService;
replicateSchemaInCluster metodu
Call stack şöyle
replicateSchemaInCluster:140, ClientSchemaService (com.hazelcast.client.impl.clientside)
put:85, ClientSchemaService (com.hazelcast.client.impl.clientside)
putToSchemaService:137, CompactStreamSerializer (com.hazelcast.internal.serialization.impl.compact)
writeObject:148, CompactStreamSerializer (com.hazelcast.internal.serialization.impl.compact)
write:116, CompactStreamSerializer (com.hazelcast.internal.serialization.impl.compact)
write:109, CompactStreamSerializer (com.hazelcast.internal.serialization.impl.compact)
write:39, StreamSerializerAdapter (com.hazelcast.internal.serialization.impl)
toBytes:238, AbstractSerializationService (com.hazelcast.internal.serialization.impl)
toBytes:217, AbstractSerializationService (com.hazelcast.internal.serialization.impl)
toData:202, AbstractSerializationService (com.hazelcast.internal.serialization.impl)
toData:157, AbstractSerializationService (com.hazelcast.internal.serialization.impl)
toData:76, ClientProxy (com.hazelcast.client.impl.spi)
putInternal:553, ClientMapProxy (com.hazelcast.client.impl.proxy)
put:275, ClientMapProxy (com.hazelcast.client.impl.proxy)
Client CompactSerializer'ın kullandığı com.hazelcast.internal.serialization.impl.compact.Schema
nesnesini sunucuya gönderir. Böylece sunucu da bu nesnesi okuyabilir.

HazelcastAPI SchemaService Arayüzü

Giriş
Şu satırı dahil ederiz
import com.hazelcast.internal.serialization.impl.compact.SchemaService;
Kalıtım şöyle
SchemaService
  MemberSchemaService
  

5 Aralık 2023 Salı

Compact Serialization Custom Configuration

Örnek
Şu satırı dahil ederiz
import com.hazelcast.nio.serialization.compact.CompactSerializer;
import com.hazelcast.nio.serialization.compact.CompactReader;
import com.hazelcast.nio.serialization.compact.CompactWriter;
Örnek - Member Side Custom Configuration
Şöyle yaparız
public class Employee {
   ...
}

public class EmployeeSerializer implements CompactSerializer<Employee> {
  @Override
  public Employee read(CompactReader reader) {
    long id = reader.readInt64("id");
    String name = reader.readString("name");
    return new Employee(id, name);
  }

  @Override
  public void write(CompactWriter writer, Employee employee) {
    writer.writeInt64("id", employee.getId());
    writer.writeString("name", employee.getName());
  }

  @Override
  public Class<Employee> getCompactClass() {
    return Employee.class;
  }

  @Override
  public String getTypeName() {
    return "employee";
  }
}

Config config = new Config();
config.getSerializationConfig()
        .getCompactSerializationConfig()
        .addSerializer(new EmployeeSerializer());
Örnek - Client Side Custom Configuration
İstemci tarafından IMap.put işlemi için call stack şöyle . Yazılan byte[] nesnesi HeapData nesnesine çevrili.
write:16, MyWorkerSerializer (com.colak.serilization.compact.serializerconfiguration)
write:7, MyWorkerSerializer (com.colak.serilization.compact.serializerconfiguration)
buildSchema:396, CompactStreamSerializer (com.hazelcast.internal.serialization.impl.compact)
writeObject:147, CompactStreamSerializer (com.hazelcast.internal.serialization.impl.compact)
write:116, CompactStreamSerializer (com.hazelcast.internal.serialization.impl.compact)
write:109, CompactStreamSerializer (com.hazelcast.internal.serialization.impl.compact)
write:39, StreamSerializerAdapter (com.hazelcast.internal.serialization.impl)
toBytes:238, AbstractSerializationService (com.hazelcast.internal.serialization.impl)
toBytes:217, AbstractSerializationService (com.hazelcast.internal.serialization.impl)
toData:202, AbstractSerializationService (com.hazelcast.internal.serialization.impl)
toData:157, AbstractSerializationService (com.hazelcast.internal.serialization.impl)
toData:76, ClientProxy (com.hazelcast.client.impl.spi)
putInternal:553, ClientMapProxy (com.hazelcast.client.impl.proxy)
put:275, ClientMapProxy (com.hazelcast.client.impl.proxy)

29 Kasım 2023 Çarşamba

ProcessorMetaSupplier.Context Arayüzü

Giriş
Şu satırı dahil ederiz
import com.hazelcast.jet.core.ProcessorMetaSupplier.Context;
Bu arayüz ProcessorMetaSupplier nesnesinin init() metoduna parametre olarak geçilir. Böylece ProcessorMetaSupplier nesnesi bazı ortam değişkenlerine erişebilir. 

Metodlar şöyle
HazelcastInstance hazelcastInstance();
JetInstance jetInstance();
long jobId();
long executionId();
JobConfig jobConfig();
int totalParallelism();
int localParallelism();
int memberCount();
String vertexName();
ILogger logger();
boolean snapshottingEnabled();
ProcessingGuarantee processingGuarantee();
long maxProcessorAccumulatedRecords();
boolean isLightJob();
Map<Address, int[]> partitionAssignment();
ClassLoader classLoader();
DataConnectionService dataConnectionService();
void checkPermission(@Nonnull Permission permission)


SchedLock

Giriş
Önemli sınıflar şöyle
1. HazelcastLockProvider
2. HazelcastLock

HazelcastLockProvider Sınıfı
lock metodu
Kod şöyle
public Optional<SimpleLock> lock(@NonNull LockConfiguration lockConfiguration) {

  final Instant now = ClockProvider.now();
  final String lockName = lockConfiguration.getName();
  final IMap<String, HazelcastLock> store = getStore();
  try {
    // lock the map key entry
    store.lock(lockName, keyLockTime(lockConfiguration), TimeUnit.MILLISECONDS);
    // just one thread at a time, in the cluster, can run this code
    // each thread waits until the lock to be unlock
    if (tryLock(lockConfiguration, now)) {
      return Optional.of(new HazelcastSimpleLock(this, lockConfiguration));
    }
  } finally {
    // released the map lock for the others threads
    store.unlock(lockName);
  }
  return Optional.empty();
}
  1. shedlock_storage isimli IMap veri yapısında HazelcastLock nesneleri saklanır. 
  2. eyLockTime() metodu now + LockAtMostUntil değerini verir. Yani IMap.lock() çağrısı gelecekteki bir zaman kadar bu key değerini kilitler
  3. tryLock() metodu kilitli olan key değerinde 
    1. HazelcastLock nesnesi yoksa yeni bir tane ekler ve true döner
    2.  HazelcastLock nesnesi varsa ve bayatlamışsa, yeni HazelcastLock  ile değiştirir ve true döner
    3. HazelcastLock nesnesi varsa ve bayatlamamışsa, kilitleyemediği için false döner





28 Kasım 2023 Salı

hazelcast.xml Cache Ayarları

Örnek
Şöyle yaparız
<cache name="*">
  <statistics-enabled>false</statistics-enabled>
  <management-enabled>false</management-enabled>
  <in-memory-format>BINARY</in-memory-format>
  <expiry-policy-factory>
    <timed-expiry-policy-factory expiry-policy-type="ETERNAL"/>
  </expiry-policy-factory>
  <eviction eviction-policy="LRU" max-size-policy="ENTRY_COUNT" size="10000"/>
</cache>

hazelcast.xml - Advanced Network Ayarları

Giriş
Network Join Ayarları yazısına bakabilirsiniz


Örnek
Şöyle yaparız
<advanced-network enabled="true">
  <join>
    <auto-detection enabled="false"/>
    <multicast enabled="true">
      <multicast-group>224.2.2.3</multicast-group>
      <multicast-port>54327</multicast-port>
      <multicast-time-to-live>32</multicast-time-to-live>
      <multicast-timeout-seconds>5</multicast-timeout-seconds>
    </multicast>
  </join>

  <member-server-socket-endpoint-config>
    <port>5701</port>
    <socket-options>
      <keep-alive>true</keep-alive>
      <tcp-no-delay>true</tcp-no-delay>
      <buffer-direct>true</buffer-direct>
    </socket-options>
  </member-server-socket-endpoint-config>
  <client-server-socket-endpoint-config>
    <port>9090</port>
    <socket-options>
    <keep-alive>true</keep-alive>
    <tcp-no-delay>true</tcp-no-delay>
    <buffer-direct>true</buffer-direct>
    </socket-options>
  </client-server-socket-endpoint-config>
</advanced-network>

21 Kasım 2023 Salı

HazelcastAPI AbstractRecordStore Sınıfı

Giriş
Şu satırı dahil ederiz
import com.hazelcast.map.impl.AbstractRecordStore;
Kodu şöyle
abstract class AbstractRecordStore implements RecordStore<Record> {
  protected final int partitionId;
  protected final String name;
  protected final LockStore lockStore;
  protected final MapContainer mapContainer;
  protected final RecordFactory recordFactory;
  protected final InMemoryFormat inMemoryFormat;
  protected final MapStoreContext mapStoreContext;
  protected final ValueComparator valueComparator;
  protected final MapServiceContext mapServiceContext;
  protected final MapDataStore<Data, Object> mapDataStore;
  protected final SerializationService serializationService;
  protected final CompositeMutationObserver<Record> mutationObserver;
  protected final LocalRecordStoreStatsImpl stats = new LocalRecordStoreStatsImpl();
  protected Storage<Data, Record> storage;
  protected IndexingMutationObserver<Record> indexingObserver;
  ...
}
storage Alanı
protected Storage<Data, Record> storage
 alanı sanırım veriyi saklayan yapı

mapDataStore Alanı
MapDataStore<Data, Object> mapDataStore alanı write-throughwrite-behind işlemlerini gerçekleştiren arayüz. WriteThroughStore yazısına bakabilirsiniz.

MapDataStore (WriteThroughStore) -> has MapStoreWrapper 
MapStoreWrapper -> has MapLoader or MapStore 

20 Kasım 2023 Pazartesi

IdentifiedDataSerializable Serialization

Giriş
Açıklaması şöyle
IdentifiedDataSerializable extends DataSerializable and introduces the following methods:
  int getClassId();
  int getFactoryId();

IdentifiedDataSerializable uses getClassId() instead of class name and it uses getFactoryId() to load the class when given the id. 

To complete the implementation, you should also implement
com.hazelcast.nio.serialization.DataSerializableFactory and register it into SerializationConfig, which can be accessed from Config.getSerializationConfig()

Factory’s responsibility is to return an instance of the right IdentifiedDataSerializable object, given the id. 

1 Kasım 2023 Çarşamba

Stream-to-Batch Join

Giriş
Join işlemi iki farklı başlık altında ele alınabilir

Stream-to-Batch Join
Kafka ve JDBC tabloları arasındaki join

Batch-to-Batch Join
JDBC tabloları arasındaki join

1. Stream-to-Batch Join
Apache Calcite bir şekilde JDBC  tarafında kullanılacak indeksleri bulamıyor. hazelcast-sql modülündeki JetJoinInfo sınıfının hem leftEquiJoinIndices hem de rightEquiJoinIndices dizisi boş geliyor. Yani JDBC tarafında hangi sütunlar için WHERE koşulu çalıştırılacağını bulamıyoruz. Bu yüzden JDBC tablosu için
SELECT column1, column2 FROM mytable;
şeklinde bir sorgu çalıştırılıyor.  
- SELECT bölümündeki çekilecek sütunlar listesi List<RexNode> projection değişkeninden geliyor. Aslında SELECT * yapılsa da olurdu
- Tüm JDBC tablosu üzerinde yürüyerek sol taraftaki satır ile sağ taraftaki satırın belirtilen koşula uyup uymadığı kontrol ediliyor. Yani aslında FullScan yapılıyor. 

2. Batch-to-Batch Join
Apache Calcite bir şekilde sağ taraftaki JDBC tablosunda kullanılacak indeksleri buluyor. JetJoinInfo sınıfının rightEquiJoinIndices dizisi dolu geliyor. Bu durumda sağ tabloda hangi sütunlar için WHERE koşulu çalıştırılacağını bulabiliyoruz. Bu yüzden JDBC tablosu için
SELECT colum1, colum2 FROM mytable WHERE colum1 = ? AND column2 = ?;
şeklinde bir sorgu çalıştırılıyor. 
- SELECT bölümündeki çekilecek sütunlar listesi List<RexNode> projection değişkeninden geliyor. 
- WHERE bölümündeki değişkenler listesi rightEquiJoinIndices değişkeninden geliyor. 
- Satırları sorgulama aşamasında yani PreparedStatement ile sorgularken soru işaretlerinin yerine de sol taraftan gelen satırdaki değerler bağlanıyor. Yani bir  anlamda IndexScan yapılıyor. 


31 Ekim 2023 Salı

Hazelcast Jet Sources.mapJournal metodu

Örnek
Şöyle yaparız
IMap<Long, String> myMap = ...;
Pipeline p = Pipeline.create();
p.readFrom(Sources.mapJournal(myMap, START_FROM_CURRENT))
  .withoutTimestamps()
  .writeTo(Sinks.jdbc("%some update query%", () -> {
    BaseDataSource dataSource = new PGXADataSource();
    dataSource.setUrl("jdbc:postgresql://localhost:5432/my_db");
    dataSource.setUser("postgres");
    dataSource.setPassword("postgres");
    dataSource.setDatabaseName("my_db");
    return dataSource;
  }, (stmt, record) -> {
    // fill query params and execute                         
}));


HazelcastAPI CP SubSystem vs Transaciton Yapılar

Giriş Bu yazının yazılma sebebi bu soru 1. Transactional Yapılar Hazelcast TransactionalMap, birden fazla veri değişikliğinin tek bir işlem ...